diff --git a/products/managed_warehouse/backend/README.md b/products/managed_warehouse/backend/README.md index 7e5d2f5be756..fb3af72bc5a0 100644 --- a/products/managed_warehouse/backend/README.md +++ b/products/managed_warehouse/backend/README.md @@ -66,10 +66,12 @@ Every copy is written to a deterministic schema inside DuckLake. Each workflow n ### Data Imports and Data Import Registration - **Schema**: `posthog_data_imports_team_` -- **Table**: `__` (prefix is user-defined on the external data source) +- **Table**: a physical name derived from the organization's naming version - **Example**: `ducklake.posthog_data_imports_team_123.stripe_prod_invoices` - **Registered files**: `s3://///_imports///` +Duckgres stores a table-naming version on the organization. Organizations that existed when versioning was introduced keep the batch sink's snake-case format, such as `tik_tok_ads_ad_report`. New organizations use the copy workflow format, such as `tiktokads_ad_report`. Copy, registration, the batch sink, and query binding derive the same physical name from that organization-level policy. Do not change the policy after an organization has written data unless the underlying tables are migrated at the same time. + Each completed import creates a timestamped prepared Parquet snapshot in the data warehouse bucket. The registration workflow copies those objects directly into the DuckLake bucket, preserving Hive partition directories, registers the destination objects with `ducklake_add_data_files`, verifies the shadow table's row count, and only then swaps it into the stable table name through the Duckgres PostgreSQL connection. Registration, verification, and the swap share one catalog transaction, so a mismatch leaves the previous table live. Each import job gets its own object prefix and child workflow ID, so a later sync does not append into the previous snapshot. The registered objects are permanent DuckLake data files, not staging files. Old generations remain reachable through DuckLake snapshots until snapshot expiration and old-file cleanup make them eligible for object deletion. Choose the bucket lifecycle policy with that retention behavior in mind. diff --git a/products/managed_warehouse/backend/common.py b/products/managed_warehouse/backend/common.py index d17e838dcf08..85e74869e05c 100644 --- a/products/managed_warehouse/backend/common.py +++ b/products/managed_warehouse/backend/common.py @@ -28,6 +28,9 @@ from psycopg import sql from tenacity import retry, retry_if_exception_type, stop_after_attempt, wait_fixed +from products.managed_warehouse.backend.facade import team_state as team_state_facade +from products.warehouse_sources.backend.facade.duckgres import duckgres_data_imports_table_name_for_version + if TYPE_CHECKING: from clickhouse_driver import Client @@ -544,17 +547,13 @@ def duckgres_data_imports_schema(team_id: int) -> str: def duckgres_data_imports_table_name(schema: ExternalDataSchema) -> str: - """Resolve the duckgres table name the data-import copy workflow writes a schema's snapshot into. + """Resolve a data-import table name from the organization's control-plane naming policy.""" - Must stay byte-identical to what the copy workflow computes so the reader resolves to the same - table the writer produced. - """ - source_type = schema.source.source_type - prefix = schema.source.prefix - normalized_name = schema.normalized_name - return sanitize_ducklake_identifier( - f"{source_type}_{prefix}_{normalized_name}" if prefix else f"{source_type}_{normalized_name}", - default_prefix="data_import", + return duckgres_data_imports_table_name_for_version( + schema.source.source_type, + schema.source.prefix, + schema.normalized_name, + team_state_facade.data_imports_table_naming_version(schema.team_id), ) diff --git a/products/managed_warehouse/backend/cp_teams.py b/products/managed_warehouse/backend/cp_teams.py index 419e52d6afc5..b38d38264778 100644 --- a/products/managed_warehouse/backend/cp_teams.py +++ b/products/managed_warehouse/backend/cp_teams.py @@ -56,6 +56,7 @@ class CPTeam: persons_table_name: str | None schema_data_imports_name: str | None earliest_event_date: date | None + data_imports_table_naming_version: str = "legacy_batch_v1" @property def resolved_events_table(self) -> str: @@ -112,6 +113,7 @@ def team_from_row(row: dict, *, organization_id: str | None = None) -> CPTeam | persons_table_name=row.get("persons_table_name") or None, schema_data_imports_name=row.get("schema_data_imports_name") or None, earliest_event_date=_parse_earliest_event_date(row.get("earliest_event_date")), + data_imports_table_naming_version=str(row.get("data_imports_table_naming_version") or "legacy_batch_v1"), ) @@ -163,11 +165,19 @@ def _rows_from_response(response: http_requests.Response) -> list[dict] | None: data = response.json() except ValueError: return None + naming_version: object = None if isinstance(data, dict): + naming_version = data.get("data_imports_table_naming_version") data = data.get("teams") if not isinstance(data, list): return None - return [row for row in data if isinstance(row, dict)] + return [ + {**row, "data_imports_table_naming_version": naming_version} + if isinstance(naming_version, str) and naming_version + else row + for row in data + if isinstance(row, dict) + ] def _fetch_rows(*, organization_id: str | None) -> list[dict] | None: diff --git a/products/managed_warehouse/backend/facade/team_state.py b/products/managed_warehouse/backend/facade/team_state.py index a3b9d4540f17..a14b6dcb8a11 100644 --- a/products/managed_warehouse/backend/facade/team_state.py +++ b/products/managed_warehouse/backend/facade/team_state.py @@ -10,6 +10,7 @@ __all__ = [ "backfill_row_exists", "data_imports_schema", + "data_imports_table_naming_version", "list_enabled_backfill_team_memberships", "resolve_events_persons_tables", "team_backfill_membership", @@ -29,6 +30,12 @@ def data_imports_schema(team_id: int) -> str: return team_state.data_imports_schema(team_id) +def data_imports_table_naming_version(team_id: int) -> str: + from products.managed_warehouse.backend import team_state + + return team_state.data_imports_table_naming_version(team_id) + + def team_backfill_membership(team_id: int) -> ManagedWarehouseTeamMembership | None: from products.managed_warehouse.backend import team_state diff --git a/products/managed_warehouse/backend/presentation/views.py b/products/managed_warehouse/backend/presentation/views.py index b8cd031167b5..e8b7bdedf0a5 100644 --- a/products/managed_warehouse/backend/presentation/views.py +++ b/products/managed_warehouse/backend/presentation/views.py @@ -657,7 +657,13 @@ def _teams_from_response(resp: Response) -> list[dict] | None: return None data = resp.data if isinstance(data, dict): + naming_version = data.get("data_imports_table_naming_version") data = data.get("teams") + if isinstance(data, list) and isinstance(naming_version, str) and naming_version: + data = [ + {**row, "data_imports_table_naming_version": naming_version} if isinstance(row, dict) else row + for row in data + ] if not isinstance(data, list): return None return [row for row in data if isinstance(row, dict)] diff --git a/products/managed_warehouse/backend/team_state.py b/products/managed_warehouse/backend/team_state.py index 1469a7e4f592..0e6463063739 100644 --- a/products/managed_warehouse/backend/team_state.py +++ b/products/managed_warehouse/backend/team_state.py @@ -94,6 +94,14 @@ def data_imports_schema(team_id: int) -> str: return schema +def data_imports_table_naming_version(team_id: int) -> str: + """The organization-level naming policy shared by Duckgres data-import readers and writers.""" + row = _get_cp_row(team_id) + if row is None: + return "copy_v1" + return row.data_imports_table_naming_version + + # --- backfill state (warehouse-status UI) ----------------------------------------- diff --git a/products/managed_warehouse/backend/temporal/ducklake_copy_data_imports_workflow.py b/products/managed_warehouse/backend/temporal/ducklake_copy_data_imports_workflow.py index 41bfb069db3a..b62a82f2cd75 100644 --- a/products/managed_warehouse/backend/temporal/ducklake_copy_data_imports_workflow.py +++ b/products/managed_warehouse/backend/temporal/ducklake_copy_data_imports_workflow.py @@ -333,7 +333,7 @@ async def _prepare_data_imports_ducklake_metadata( source_normalized_name=normalized_name, source_table_uri=source_table_uri, ducklake_schema_name=ducklake_schema_name, - ducklake_table_name=duckgres_data_imports_table_name(schema), + ducklake_table_name=await database_sync_to_async(duckgres_data_imports_table_name)(schema), verification_queries=list(get_data_imports_verification_queries(normalized_name)), source_partition_column=partition_column, staging_uri=staging_uri, diff --git a/products/managed_warehouse/backend/temporal/ducklake_register_data_imports_workflow.py b/products/managed_warehouse/backend/temporal/ducklake_register_data_imports_workflow.py index 40e41ee0ca96..d110c5ab787f 100644 --- a/products/managed_warehouse/backend/temporal/ducklake_register_data_imports_workflow.py +++ b/products/managed_warehouse/backend/temporal/ducklake_register_data_imports_workflow.py @@ -184,7 +184,7 @@ async def prepare_ducklake_data_imports_registration_activity( prepared_source_uri = f"{settings.BUCKET_URL}/{schema.folder_path()}/{inputs.prepared_queryable_folder}" ducklake_schema_name = await database_sync_to_async(duckgres_data_imports_schema)(inputs.team_id) - ducklake_table_name = duckgres_data_imports_table_name(schema) + ducklake_table_name = await database_sync_to_async(duckgres_data_imports_table_name)(schema) landing_uri = await database_sync_to_async(_resolve_data_imports_landing_uri)( team_id=inputs.team_id, ducklake_schema_name=ducklake_schema_name, diff --git a/products/managed_warehouse/backend/tests/api/test_presentation.py b/products/managed_warehouse/backend/tests/api/test_presentation.py index 20294766ca39..f611a1cd63d1 100644 --- a/products/managed_warehouse/backend/tests/api/test_presentation.py +++ b/products/managed_warehouse/backend/tests/api/test_presentation.py @@ -492,6 +492,21 @@ def test_cp_bucket_for_returns_none_when_cp_has_no_bucket(mock_request: MagicMoc assert managed_warehouse.cp_bucket_for(org.id) is None +def test_teams_response_attaches_the_org_naming_policy_to_each_row() -> None: + response = Response( + { + "teams": [{"team_id": 1}, {"team_id": 2}], + "data_imports_table_naming_version": "copy_v1", + }, + status=200, + ) + + assert managed_warehouse._teams_from_response(response) == [ + {"team_id": 1, "data_imports_table_naming_version": "copy_v1"}, + {"team_id": 2, "data_imports_table_naming_version": "copy_v1"}, + ] + + @pytest.mark.django_db @patch("products.managed_warehouse.backend.presentation.views.is_enabled", return_value=True) @patch("products.managed_warehouse.backend.presentation.views._request") diff --git a/products/managed_warehouse/backend/tests/temporal/test_ducklake_copy_data_imports_workflow.py b/products/managed_warehouse/backend/tests/temporal/test_ducklake_copy_data_imports_workflow.py index e2e605fe1709..29ee5e8d355e 100644 --- a/products/managed_warehouse/backend/tests/temporal/test_ducklake_copy_data_imports_workflow.py +++ b/products/managed_warehouse/backend/tests/temporal/test_ducklake_copy_data_imports_workflow.py @@ -3,7 +3,7 @@ from collections.abc import Sequence import pytest -from unittest.mock import AsyncMock, MagicMock +from unittest.mock import AsyncMock, MagicMock, patch from django.conf import settings from django.test import override_settings @@ -45,12 +45,15 @@ @pytest.fixture(autouse=True) def _cp_no_rows(): - # The workflow reads the schema through the typed team-state facade. - from unittest.mock import patch - - with patch( - "products.managed_warehouse.backend.facade.team_state.data_imports_schema", - side_effect=lambda team_id: f"posthog_data_imports_team_{team_id}", + with ( + patch( + "products.managed_warehouse.backend.facade.team_state.data_imports_schema", + side_effect=lambda team_id: f"posthog_data_imports_team_{team_id}", + ), + patch( + "products.managed_warehouse.backend.facade.team_state.data_imports_table_naming_version", + return_value="copy_v1", + ), ): yield diff --git a/products/managed_warehouse/backend/tests/temporal/test_ducklake_register_data_imports_workflow.py b/products/managed_warehouse/backend/tests/temporal/test_ducklake_register_data_imports_workflow.py index feabd20abbce..ada5c74ae099 100644 --- a/products/managed_warehouse/backend/tests/temporal/test_ducklake_register_data_imports_workflow.py +++ b/products/managed_warehouse/backend/tests/temporal/test_ducklake_register_data_imports_workflow.py @@ -3,7 +3,7 @@ import contextlib import pytest -from unittest.mock import AsyncMock, MagicMock +from unittest.mock import AsyncMock, MagicMock, patch from parameterized import parameterized from temporalio.exceptions import ApplicationError @@ -33,11 +33,15 @@ @pytest.fixture(autouse=True) def _cp_no_rows(): - from unittest.mock import patch - - with patch( - "products.managed_warehouse.backend.facade.team_state.data_imports_schema", - side_effect=lambda team_id: f"posthog_data_imports_team_{team_id}", + with ( + patch( + "products.managed_warehouse.backend.facade.team_state.data_imports_schema", + side_effect=lambda team_id: f"posthog_data_imports_team_{team_id}", + ), + patch( + "products.managed_warehouse.backend.facade.team_state.data_imports_table_naming_version", + return_value="copy_v1", + ), ): yield diff --git a/products/managed_warehouse/backend/tests/test_client.py b/products/managed_warehouse/backend/tests/test_client.py index e720ca21fdb0..94fe9eea86f0 100644 --- a/products/managed_warehouse/backend/tests/test_client.py +++ b/products/managed_warehouse/backend/tests/test_client.py @@ -16,11 +16,17 @@ @pytest.fixture(autouse=True) def _cp_no_rows(): - # Compilation resolves the data-import schema through the typed team-state facade; - # keep these tests independent from the control plane. - with mock.patch( - "products.managed_warehouse.backend.facade.team_state.data_imports_schema", - side_effect=lambda team_id: f"posthog_data_imports_team_{team_id}", + # Compilation resolves managed table metadata through the typed team-state facade, + # so pin both organization policies to keep these tests independent from the control plane. + with ( + mock.patch( + "products.managed_warehouse.backend.facade.team_state.data_imports_schema", + side_effect=lambda team_id: f"posthog_data_imports_team_{team_id}", + ), + mock.patch( + "products.managed_warehouse.backend.facade.team_state.data_imports_table_naming_version", + return_value="copy_v1", + ), ): yield diff --git a/products/managed_warehouse/backend/tests/test_common.py b/products/managed_warehouse/backend/tests/test_common.py index 40c3aad15ca0..fb230998e3b2 100644 --- a/products/managed_warehouse/backend/tests/test_common.py +++ b/products/managed_warehouse/backend/tests/test_common.py @@ -10,6 +10,7 @@ from products.managed_warehouse.backend.common import ( default_bucket_region, + duckgres_data_imports_table_name, initialize_ducklake, is_version_mismatch, reset_ducklake_catalog, @@ -54,6 +55,36 @@ def test_region_follows_cloud_deployment(self, _name, deployment, expected): assert default_bucket_region() == expected +class TestDuckgresDataImportsTableName: + @parameterized.expand( + [ + ("copy_mysql", "copy_v1", "MySQL", "SalesEU", "orders", "mysql_saleseu_orders"), + ("copy_google_ads", "copy_v1", "GoogleAds", None, "video", "googleads_video"), + ("legacy_batch_tiktok", "legacy_batch_v1", "TikTokAds", None, "video", "tik_tok_ads_video"), + ] + ) + def test_uses_the_org_policy( + self, + _name: str, + naming_version: str, + source_type: str, + prefix: str | None, + normalized_name: str, + expected: str, + ) -> None: + schema = MagicMock() + schema.source.source_type = source_type + schema.source.prefix = prefix + schema.normalized_name = normalized_name + schema.team_id = 1 + + with patch( + "products.managed_warehouse.backend.team_state.data_imports_table_naming_version", + return_value=naming_version, + ): + assert duckgres_data_imports_table_name(schema) == expected + + TEST_CONFIG = { "DUCKLAKE_RDS_HOST": "localhost", "DUCKLAKE_RDS_PORT": "5432", diff --git a/products/managed_warehouse/backend/tests/test_cp_teams.py b/products/managed_warehouse/backend/tests/test_cp_teams.py index 85c461b3d04c..8de7afb8ef25 100644 --- a/products/managed_warehouse/backend/tests/test_cp_teams.py +++ b/products/managed_warehouse/backend/tests/test_cp_teams.py @@ -22,6 +22,7 @@ def _row(**overrides) -> dict: "persons_table_name": None, "schema_data_imports_name": None, "earliest_event_date": None, + "data_imports_table_naming_version": "copy_v1", } row.update(overrides) return row @@ -81,8 +82,18 @@ def test_coerces_types_defensively(self) -> None: persons_table_name=None, schema_data_imports_name=None, earliest_event_date=date(2020, 6, 15), + data_imports_table_naming_version="copy_v1", ) + def test_missing_naming_version_defaults_to_legacy_batch_during_rolling_deploy(self) -> None: + row = _row() + row.pop("data_imports_table_naming_version") + + team = team_from_row(row) + + assert team is not None + assert team.data_imports_table_naming_version == "legacy_batch_v1" + @parameterized.expand( [ ("missing_team_id", {"team_id": None}), @@ -203,8 +214,13 @@ def test_list_enabled_backfills_filters_disabled_rows(self) -> None: class TestControlPlaneTransport: @override_settings(DUCKGRES_API_URL="https://duckgres.example/", DUCKGRES_INTERNAL_SECRET="secret") def test_org_read_uses_internal_transport_with_the_control_plane_request_shape(self) -> None: + response_row = _row() + response_row.pop("data_imports_table_naming_version") response = MagicMock(status_code=200, text="ok") - response.json.return_value = {"teams": [_row()]} + response.json.return_value = { + "teams": [response_row], + "data_imports_table_naming_version": "copy_v1", + } with patch( "products.managed_warehouse.backend.cp_teams.internal_requests.request", diff --git a/products/managed_warehouse/backend/tests/test_team_state.py b/products/managed_warehouse/backend/tests/test_team_state.py index 58a184b46dd4..340abd924441 100644 --- a/products/managed_warehouse/backend/tests/test_team_state.py +++ b/products/managed_warehouse/backend/tests/test_team_state.py @@ -39,6 +39,7 @@ def _cp_row(team: Team, schema_name: str, **overrides) -> dict: "persons_table_name": None, "schema_data_imports_name": None, "earliest_event_date": None, + "data_imports_table_naming_version": "copy_v1", } row.update(overrides) return row @@ -74,6 +75,19 @@ def test_serves_cached_rows_during_an_outage(self) -> None: assert team_state.data_imports_schema(team.id) == "posthog_data_imports_cp_schema" +@pytest.mark.django_db +class TestDataImportsTableNamingVersion: + def test_resolves_the_org_policy_from_the_cp_row(self) -> None: + org, team = _team() + with _patch_org_rows([_cp_row(team, "cp_schema", data_imports_table_naming_version="legacy_batch_v1")]): + assert team_state.data_imports_table_naming_version(team.id) == "legacy_batch_v1" + + def test_without_a_cp_row_uses_copy_workflow_naming(self) -> None: + org, team = _team() + with _patch_org_rows([]): + assert team_state.data_imports_table_naming_version(team.id) == "copy_v1" + + @pytest.mark.django_db class TestEventsPersonsTables: @parameterized.expand( diff --git a/products/warehouse_sources/backend/duckgres_naming.py b/products/warehouse_sources/backend/duckgres_naming.py new file mode 100644 index 000000000000..cbbf904eb528 --- /dev/null +++ b/products/warehouse_sources/backend/duckgres_naming.py @@ -0,0 +1,31 @@ +from __future__ import annotations + +import re + +from products.warehouse_sources.backend.temporal.data_imports.naming_convention import NamingConvention + +_IDENTIFIER_SANITIZE_RE = re.compile(r"[^0-9a-zA-Z]+") +_DUCKGRES_IDENTIFIER_MAX_LENGTH = 63 + + +def duckgres_data_imports_table_name_for_version( + source_type: str, + prefix: str | None, + normalized_name: str, + naming_version: str, +) -> str: + raw_name = f"{source_type}_{prefix}_{normalized_name}" if prefix else f"{source_type}_{normalized_name}" + if naming_version == "copy_v1": + return _sanitize_identifier(raw_name, default_prefix="data_import") + if naming_version == "legacy_batch_v1": + return NamingConvention.normalize_identifier(raw_name, max_length=_DUCKGRES_IDENTIFIER_MAX_LENGTH) + raise ValueError(f"Unsupported Duckgres data imports table naming version: {naming_version!r}") + + +def _sanitize_identifier(raw: str, *, default_prefix: str) -> str: + cleaned = _IDENTIFIER_SANITIZE_RE.sub("_", (raw or "").strip()).strip("_").lower() + if not cleaned: + cleaned = default_prefix + if cleaned[0].isdigit(): + cleaned = f"{default_prefix}_{cleaned}" + return cleaned[:_DUCKGRES_IDENTIFIER_MAX_LENGTH] diff --git a/products/warehouse_sources/backend/facade/duckgres.py b/products/warehouse_sources/backend/facade/duckgres.py new file mode 100644 index 000000000000..91d4b4cfc485 --- /dev/null +++ b/products/warehouse_sources/backend/facade/duckgres.py @@ -0,0 +1,14 @@ +from __future__ import annotations + +from products.warehouse_sources.backend.duckgres_naming import ( + duckgres_data_imports_table_name_for_version as _duckgres_data_imports_table_name_for_version, +) + + +def duckgres_data_imports_table_name_for_version( + source_type: str, + prefix: str | None, + normalized_name: str, + naming_version: str, +) -> str: + return _duckgres_data_imports_table_name_for_version(source_type, prefix, normalized_name, naming_version) diff --git a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/processor.py b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/processor.py index 446dc971c4a3..6261b0c71c39 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/processor.py +++ b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/processor.py @@ -22,6 +22,7 @@ from products.managed_warehouse.backend.facade.api import ( duckgres_data_imports_schema, + duckgres_data_imports_table_name, get_duckgres_query_server_config, setup_duckgres_session, ) @@ -604,14 +605,7 @@ def _duckgres_schema_name(team_id: int) -> str: def _duckgres_table_name(schema: ExternalDataSchema) -> str: - source_type = schema.source.source_type - normalized_name = schema.normalized_name - raw_name = ( - f"{source_type}_{schema.source.prefix}_{normalized_name}" - if schema.source.prefix - else f"{source_type}_{normalized_name}" - ) - return NamingConvention.normalize_identifier(raw_name, max_length=63) + return duckgres_data_imports_table_name(schema) def _should_replace_table(batch: PendingBatch) -> bool: diff --git a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/test_processor.py b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/test_processor.py index 46cea393cb99..ef1ed07e0e24 100644 --- a/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/test_processor.py +++ b/products/warehouse_sources/backend/temporal/data_imports/pipelines/pipeline_v3/duckgres/test_processor.py @@ -23,10 +23,15 @@ @pytest.fixture(autouse=True) def _stub_schema_resolver(): - # These no-DB unit tests pin the control-plane-backed schema resolver. - with patch( - "products.warehouse_sources.backend.temporal.data_imports.pipelines.pipeline_v3.duckgres.processor.duckgres_data_imports_schema", - return_value="posthog_data_imports_team_1", + with ( + patch( + "products.warehouse_sources.backend.temporal.data_imports.pipelines.pipeline_v3.duckgres.processor.duckgres_data_imports_schema", + return_value="posthog_data_imports_team_1", + ), + patch( + "products.warehouse_sources.backend.temporal.data_imports.pipelines.pipeline_v3.duckgres.processor.duckgres_data_imports_table_name", + return_value="stripe_customers", + ), ): yield @@ -60,6 +65,7 @@ def _make_batch(**overrides: Any) -> PendingBatch: def _make_schema() -> Mock: schema = Mock() + schema.team_id = 1 schema.normalized_name = "customers" schema.source.source_type = "Stripe" schema.source.prefix = None diff --git a/products/warehouse_sources/backend/tests/test_duckgres_naming.py b/products/warehouse_sources/backend/tests/test_duckgres_naming.py new file mode 100644 index 000000000000..8b016be3f48d --- /dev/null +++ b/products/warehouse_sources/backend/tests/test_duckgres_naming.py @@ -0,0 +1,39 @@ +import pytest + +from parameterized import parameterized + +from products.warehouse_sources.backend.duckgres_naming import duckgres_data_imports_table_name_for_version + + +class TestDuckgresDataImportsTableNameForVersion: + @parameterized.expand( + [ + ("mysql", "MySQL", "SalesEU", "customer_orders", "mysql_saleseu_customer_orders"), + ("bigquery", "BigQuery", None, "daily_stats", "bigquery_daily_stats"), + ("google_ads", "GoogleAds", None, "video", "googleads_video"), + ("tiktok_ads", "TikTokAds", "prod__us", "ad_report", "tiktokads_prod_us_ad_report"), + ] + ) + def test_source_keys_are_stable_without_camel_case_splitting( + self, _name: str, source_type: str, prefix: str | None, schema_name: str, expected: str + ) -> None: + assert duckgres_data_imports_table_name_for_version(source_type, prefix, schema_name, "copy_v1") == expected + + def test_legacy_batch_version_preserves_aggressive_snake_case(self) -> None: + assert ( + duckgres_data_imports_table_name_for_version("TikTokAds", None, "ad_report", "legacy_batch_v1") + == "tik_tok_ads_ad_report" + ) + + def test_long_names_have_stable_collision_resistant_suffixes(self) -> None: + first = duckgres_data_imports_table_name_for_version("Postgres", None, "a" * 90, "legacy_batch_v1") + second = duckgres_data_imports_table_name_for_version("Postgres", None, "a" * 89 + "b", "legacy_batch_v1") + + assert len(first) == 63 + assert len(second) == 63 + assert first != second + assert first == duckgres_data_imports_table_name_for_version("Postgres", None, "a" * 90, "legacy_batch_v1") + + def test_rejects_unknown_version(self) -> None: + with pytest.raises(ValueError, match="Unsupported Duckgres data imports table naming version"): + duckgres_data_imports_table_name_for_version("Postgres", None, "orders", "future")