From 8c4e9ea46a43d8a04736cd1f6df9b80a8acec526 Mon Sep 17 00:00:00 2001 From: eric Date: Sun, 2 Aug 2026 17:16:14 -0700 Subject: [PATCH 1/3] chore(data-warehouse): resolve post-merge conflicts --- products/managed_warehouse/backend/README.md | 4 +- products/managed_warehouse/backend/common.py | 19 ++++--- .../managed_warehouse/backend/cp_teams.py | 2 + .../backend/facade/team_state.py | 7 +++ .../backend/presentation/views.py | 6 +++ .../managed_warehouse/backend/team_state.py | 8 +++ .../ducklake_copy_data_imports_workflow.py | 4 +- ...ducklake_register_data_imports_workflow.py | 4 +- .../backend/tests/api/test_presentation.py | 15 ++++++ .../backend/tests/test_common.py | 47 ++++++++++++++++ .../backend/tests/test_cp_teams.py | 11 ++++ .../backend/tests/test_team_state.py | 14 +++++ .../backend/duckgres_naming.py | 31 +++++++++++ .../backend/duckgres_table_binding.py | 32 +++++++++++ .../backend/facade/duckgres.py | 23 ++++++++ .../0116_add_duckgres_table_name.py | 15 ++++++ .../backend/models/external_data_schema.py | 1 + .../pipeline_v3/duckgres/processor.py | 13 ++--- .../pipeline_v3/duckgres/test_processor.py | 2 + .../backend/tests/test_duckgres_naming.py | 39 ++++++++++++++ .../backend/tests/test_models.py | 54 +++++++++++++++++-- 21 files changed, 328 insertions(+), 23 deletions(-) create mode 100644 products/warehouse_sources/backend/duckgres_naming.py create mode 100644 products/warehouse_sources/backend/duckgres_table_binding.py create mode 100644 products/warehouse_sources/backend/facade/duckgres.py create mode 100644 products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py create mode 100644 products/warehouse_sources/backend/tests/test_duckgres_naming.py diff --git a/products/managed_warehouse/backend/README.md b/products/managed_warehouse/backend/README.md index 7e5d2f5be756..9a6e44b3f2bc 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 pinned on each imported schema - **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`. The first Duckgres writer atomically pins the resulting physical name on each imported schema. Copy, registration, the batch sink, and query binding then reuse that pin, so later policy changes cannot rename a table that already contains data. + 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..5bf0ee6e4970 100644 --- a/products/managed_warehouse/backend/common.py +++ b/products/managed_warehouse/backend/common.py @@ -28,6 +28,8 @@ from psycopg import sql from tenacity import retry, retry_if_exception_type, stop_after_attempt, wait_fixed +from products.warehouse_sources.backend.facade.duckgres import duckgres_data_imports_table_name_for_version + if TYPE_CHECKING: from clickhouse_driver import Client @@ -549,12 +551,17 @@ def duckgres_data_imports_table_name(schema: ExternalDataSchema) -> str: 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", + pinned_name = getattr(schema, "duckgres_table_name", None) + if isinstance(pinned_name, str) and pinned_name: + return pinned_name + + from products.managed_warehouse.backend import team_state # noqa: PLC0415 + + return duckgres_data_imports_table_name_for_version( + schema.source.source_type, + schema.source.prefix, + schema.normalized_name, + team_state.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..fc6e06438545 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"), ) 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..f2fcb60fecb4 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-stable naming version used when a Duckgres writer first binds a table.""" + 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..393bcf70827c 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 @@ -29,7 +29,6 @@ _get_org_id_for_team, attach_catalog, duckgres_data_imports_schema, - duckgres_data_imports_table_name, get_config, get_duckgres_server_by_team_org, get_duckgres_server_for_organization, @@ -64,6 +63,7 @@ get_ducklake_copy_data_imports_verification_metric, record_ducklake_copy_data_imports_stage_duration, ) +from products.warehouse_sources.backend.facade.duckgres import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.facade.models import ExternalDataSchema from products.warehouse_sources.backend.facade.pipelines import DUCKGRES_BATCH_SINK_FLAG, is_duckgres_sink_team_member @@ -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(bind_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..56e43b68ed92 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 @@ -32,7 +32,6 @@ from products.managed_warehouse.backend.common import ( _get_org_id_for_team, duckgres_data_imports_schema, - duckgres_data_imports_table_name, get_config, get_duckgres_server_by_team_org, get_duckgres_server_for_organization, @@ -50,6 +49,7 @@ get_ducklake_register_data_imports_started_metric, record_ducklake_register_data_imports_stage_duration, ) +from products.warehouse_sources.backend.facade.duckgres import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.facade.models import ExternalDataSchema LOGGER = get_logger(__name__) @@ -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(bind_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/test_common.py b/products/managed_warehouse/backend/tests/test_common.py index 40c3aad15ca0..4ebc6b93377a 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,52 @@ def test_region_follows_cloud_deployment(self, _name, deployment, expected): assert default_bucket_region() == expected +class TestDuckgresDataImportsTableName: + @parameterized.expand( + [ + ("copy_mysql", None, "copy_v1", "MySQL", "SalesEU", "orders", "mysql_saleseu_orders"), + ("copy_google_ads", None, "copy_v1", "GoogleAds", None, "video", "googleads_video"), + ("legacy_batch_tiktok", None, "legacy_batch_v1", "TikTokAds", None, "video", "tik_tok_ads_video"), + ( + "pinned_name", + "googleads_video_4f12abcd", + "legacy_batch_v1", + "GoogleAds", + None, + "video", + "googleads_video_4f12abcd", + ), + ] + ) + def test_pinned_name_wins_and_null_uses_the_org_policy( + self, + _name: str, + pinned_name: str | None, + naming_version: str, + source_type: str, + prefix: str | None, + normalized_name: str, + expected: str, + ) -> None: + schema = MagicMock() + schema.duckgres_table_name = pinned_name + 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, + ) as mock_version: + assert duckgres_data_imports_table_name(schema) == expected + + if pinned_name: + mock_version.assert_not_called() + else: + mock_version.assert_called_once_with(1) + + 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..7b0650bf01f1 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}), 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..b66be6d3630b --- /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"[^A-Za-z0-9_]+") +_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/duckgres_table_binding.py b/products/warehouse_sources/backend/duckgres_table_binding.py new file mode 100644 index 000000000000..8a15877657de --- /dev/null +++ b/products/warehouse_sources/backend/duckgres_table_binding.py @@ -0,0 +1,32 @@ +from __future__ import annotations + +from django.db.models import Q + +from products.managed_warehouse.backend.facade import team_state +from products.warehouse_sources.backend.duckgres_naming import duckgres_data_imports_table_name_for_version +from products.warehouse_sources.backend.models.external_data_schema import ExternalDataSchema + + +def bind_duckgres_data_imports_table_name(schema: ExternalDataSchema) -> str: + pinned_name = schema.duckgres_table_name + if pinned_name: + return pinned_name + + naming_version = team_state.data_imports_table_naming_version(schema.team_id) + candidate = duckgres_data_imports_table_name_for_version( + schema.source.source_type, + schema.source.prefix, + schema.normalized_name, + naming_version, + ) + updated = ExternalDataSchema.objects.filter( + Q(duckgres_table_name__isnull=True) | Q(duckgres_table_name=""), pk=schema.pk + ).update(duckgres_table_name=candidate) + if updated: + schema.duckgres_table_name = candidate + return candidate + + schema.refresh_from_db(fields=["duckgres_table_name"]) + if not schema.duckgres_table_name: + raise RuntimeError(f"External data schema {schema.pk} has no Duckgres table name after binding") + return schema.duckgres_table_name diff --git a/products/warehouse_sources/backend/facade/duckgres.py b/products/warehouse_sources/backend/facade/duckgres.py new file mode 100644 index 000000000000..33a5bc7c9b2d --- /dev/null +++ b/products/warehouse_sources/backend/facade/duckgres.py @@ -0,0 +1,23 @@ +from __future__ import annotations + +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from products.warehouse_sources.backend.facade.models import ExternalDataSchema + + +def bind_duckgres_data_imports_table_name(schema: ExternalDataSchema) -> str: + from products.warehouse_sources.backend.duckgres_table_binding import bind_duckgres_data_imports_table_name + + return bind_duckgres_data_imports_table_name(schema) + + +def duckgres_data_imports_table_name_for_version( + source_type: str, + prefix: str | None, + normalized_name: str, + naming_version: str, +) -> str: + from products.warehouse_sources.backend.duckgres_naming import duckgres_data_imports_table_name_for_version + + return duckgres_data_imports_table_name_for_version(source_type, prefix, normalized_name, naming_version) diff --git a/products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py b/products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py new file mode 100644 index 000000000000..4e00788abcb3 --- /dev/null +++ b/products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py @@ -0,0 +1,15 @@ +from django.db import migrations, models + + +class Migration(migrations.Migration): + dependencies = [ + ("warehouse_sources", "0115_scaffold_four_requested_sources"), + ] + + operations = [ + migrations.AddField( + model_name="externaldataschema", + name="duckgres_table_name", + field=models.CharField(blank=True, max_length=63, null=True), + ), + ] diff --git a/products/warehouse_sources/backend/models/external_data_schema.py b/products/warehouse_sources/backend/models/external_data_schema.py index 4729c67e6f9f..5cf93083df65 100644 --- a/products/warehouse_sources/backend/models/external_data_schema.py +++ b/products/warehouse_sources/backend/models/external_data_schema.py @@ -86,6 +86,7 @@ class SyncFrequency(models.TextChoices): # during multi-schema migration) to their original path. Empty for rows written before this # column existed — readers fall back to the legacy JSON key, then the normalized schema `name`. s3_folder_name = models.CharField(max_length=400, null=True, blank=True) + duckgres_table_name = models.CharField(max_length=63, null=True, blank=True) # Deprecated in favour of `sync_frequency_interval` sync_frequency = deprecate_field( models.CharField(max_length=128, choices=SyncFrequency, default=SyncFrequency.DAILY, blank=True) 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..0871380a16f8 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,9 +22,11 @@ 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, ) +from products.warehouse_sources.backend.duckgres_table_binding import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.models import ExternalDataJob, ExternalDataSchema from products.warehouse_sources.backend.temporal.data_imports.naming_convention import NamingConvention from products.warehouse_sources.backend.temporal.data_imports.pipelines.pipeline_v3.batch_consumer import ( @@ -244,6 +246,8 @@ def process_batch(batch: PendingBatch) -> None: raise ValueError(f"ExternalDataJob {batch.job_id} has no schema") schema = job.schema + bind_duckgres_data_imports_table_name(schema) + kind = "backfill" if _is_backfill_batch(batch) else "live" # One ORM lookup serves both the cache key and the connection config; it is # per-batch on purpose so the key always reflects the team's CURRENT org. @@ -604,14 +608,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..257bf7c53613 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 @@ -60,7 +60,9 @@ def _make_batch(**overrides: Any) -> PendingBatch: def _make_schema() -> Mock: schema = Mock() + schema.team_id = 1 schema.normalized_name = "customers" + schema.duckgres_table_name = "stripe_customers" schema.source.source_type = "Stripe" schema.source.prefix = None return schema 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..d5c35d0614e0 --- /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", "ad_report", "tiktokads_prod_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") diff --git a/products/warehouse_sources/backend/tests/test_models.py b/products/warehouse_sources/backend/tests/test_models.py index 5cfc8f799136..951f2d3058c7 100644 --- a/products/warehouse_sources/backend/tests/test_models.py +++ b/products/warehouse_sources/backend/tests/test_models.py @@ -16,6 +16,7 @@ from posthog.models.signals import model_activity_signal +from products.warehouse_sources.backend.duckgres_table_binding import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.models.credential import DataWarehouseCredential from products.warehouse_sources.backend.models.external_data_job import ExternalDataJob from products.warehouse_sources.backend.models.external_data_schema import ( @@ -76,17 +77,20 @@ def test_resolved_s3_folder_name( class TestExternalDataSchemaSave(BaseTest): - def _source(self) -> ExternalDataSource: + def _source(self, source_type: str = "Postgres", prefix: str | None = None) -> ExternalDataSource: return ExternalDataSource.objects.create( team_id=self.team.pk, source_id=str(uuid.uuid4()), connection_id=str(uuid.uuid4()), status="Completed", - source_type="Postgres", + source_type=source_type, + prefix=prefix, ) - def _create(self, name: str, **kwargs) -> ExternalDataSchema: - return ExternalDataSchema.objects.create(team_id=self.team.pk, source=self._source(), name=name, **kwargs) + def _create(self, name: str, *, source: ExternalDataSource | None = None, **kwargs) -> ExternalDataSchema: + return ExternalDataSchema.objects.create( + team_id=self.team.pk, source=source or self._source(), name=name, **kwargs + ) def test_save_populates_s3_folder_name_from_name(self) -> None: # The folder is the normalized name — never NULL for a new row. @@ -115,6 +119,48 @@ def test_partial_update_backfills_null_folder(self) -> None: schema.refresh_from_db() assert schema.s3_folder_name == "orders" + def test_new_schema_waits_for_a_duckgres_writer_to_pin_its_name(self) -> None: + schema = self._create("customer_orders", source=self._source("MySQL", "SalesEU")) + assert schema.duckgres_table_name is None + + @parameterized.expand( + [ + ("copy", "copy_v1", "tiktokads_ad_report"), + ("legacy_batch", "legacy_batch_v1", "tik_tok_ads_ad_report"), + ] + ) + def test_duckgres_writer_binds_the_org_naming_policy(self, _name: str, naming_version: str, expected: str) -> None: + schema = self._create("ad_report", source=self._source("TikTokAds")) + + with patch( + "products.warehouse_sources.backend.duckgres_table_binding.team_state.data_imports_table_naming_version", + return_value=naming_version, + ): + assert bind_duckgres_data_imports_table_name(schema) == expected + + schema.refresh_from_db() + assert schema.duckgres_table_name == expected + + def test_duckgres_writer_keeps_the_first_name_when_another_writer_wins(self) -> None: + schema = self._create("ad_report", source=self._source("TikTokAds")) + ExternalDataSchema.objects.filter(pk=schema.pk).update(duckgres_table_name="already_bound") + + with patch( + "products.warehouse_sources.backend.duckgres_table_binding.team_state.data_imports_table_naming_version", + return_value="copy_v1", + ): + assert bind_duckgres_data_imports_table_name(schema) == "already_bound" + + assert schema.duckgres_table_name == "already_bound" + + def test_existing_schema_without_pin_stays_on_legacy_naming(self) -> None: + schema = self._create("orders", source=self._source("MySQL", "SalesEU")) + schema.status = "Completed" + schema.save(update_fields=["status", "updated_at"]) + schema.refresh_from_db() + + assert schema.duckgres_table_name is None + class TestExternalDataSchemaActivityLogging(BaseTest): """Internal pipeline-driven bookkeeping saves must bypass ModelActivityMixin so they neither From effbb9507c1e35b188c18848ee76c09792afc16c Mon Sep 17 00:00:00 2001 From: eric Date: Sun, 2 Aug 2026 17:52:28 -0700 Subject: [PATCH 2/3] refactor(data-warehouse): derive Duckgres names from org policy --- products/managed_warehouse/backend/README.md | 4 +- products/managed_warehouse/backend/common.py | 14 ++--- .../managed_warehouse/backend/cp_teams.py | 10 +++- .../managed_warehouse/backend/team_state.py | 2 +- .../ducklake_copy_data_imports_workflow.py | 4 +- ...ducklake_register_data_imports_workflow.py | 4 +- ...est_ducklake_copy_data_imports_workflow.py | 17 +++--- ...ducklake_register_data_imports_workflow.py | 16 +++--- .../backend/tests/test_common.py | 26 ++------- .../backend/tests/test_cp_teams.py | 7 ++- .../backend/duckgres_naming.py | 2 +- .../backend/duckgres_table_binding.py | 32 ----------- .../backend/facade/duckgres.py | 17 ++---- .../0116_add_duckgres_table_name.py | 15 ------ .../backend/models/external_data_schema.py | 1 - .../pipeline_v3/duckgres/processor.py | 3 -- .../pipeline_v3/duckgres/test_processor.py | 14 +++-- .../backend/tests/test_duckgres_naming.py | 2 +- .../backend/tests/test_models.py | 54 ++----------------- 19 files changed, 69 insertions(+), 175 deletions(-) delete mode 100644 products/warehouse_sources/backend/duckgres_table_binding.py delete mode 100644 products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py diff --git a/products/managed_warehouse/backend/README.md b/products/managed_warehouse/backend/README.md index 9a6e44b3f2bc..fb3af72bc5a0 100644 --- a/products/managed_warehouse/backend/README.md +++ b/products/managed_warehouse/backend/README.md @@ -66,11 +66,11 @@ 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**: a physical name pinned on each imported schema +- **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`. The first Duckgres writer atomically pins the resulting physical name on each imported schema. Copy, registration, the batch sink, and query binding then reuse that pin, so later policy changes cannot rename a table that already contains data. +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. diff --git a/products/managed_warehouse/backend/common.py b/products/managed_warehouse/backend/common.py index 5bf0ee6e4970..85e74869e05c 100644 --- a/products/managed_warehouse/backend/common.py +++ b/products/managed_warehouse/backend/common.py @@ -28,6 +28,7 @@ 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: @@ -546,22 +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. - - Must stay byte-identical to what the copy workflow computes so the reader resolves to the same - table the writer produced. - """ - pinned_name = getattr(schema, "duckgres_table_name", None) - if isinstance(pinned_name, str) and pinned_name: - return pinned_name - - from products.managed_warehouse.backend import team_state # noqa: PLC0415 + """Resolve a data-import table name from the organization's control-plane naming policy.""" return duckgres_data_imports_table_name_for_version( schema.source.source_type, schema.source.prefix, schema.normalized_name, - team_state.data_imports_table_naming_version(schema.team_id), + 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 fc6e06438545..b38d38264778 100644 --- a/products/managed_warehouse/backend/cp_teams.py +++ b/products/managed_warehouse/backend/cp_teams.py @@ -165,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/team_state.py b/products/managed_warehouse/backend/team_state.py index f2fcb60fecb4..0e6463063739 100644 --- a/products/managed_warehouse/backend/team_state.py +++ b/products/managed_warehouse/backend/team_state.py @@ -95,7 +95,7 @@ def data_imports_schema(team_id: int) -> str: def data_imports_table_naming_version(team_id: int) -> str: - """The organization-stable naming version used when a Duckgres writer first binds a table.""" + """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" 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 393bcf70827c..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 @@ -29,6 +29,7 @@ _get_org_id_for_team, attach_catalog, duckgres_data_imports_schema, + duckgres_data_imports_table_name, get_config, get_duckgres_server_by_team_org, get_duckgres_server_for_organization, @@ -63,7 +64,6 @@ get_ducklake_copy_data_imports_verification_metric, record_ducklake_copy_data_imports_stage_duration, ) -from products.warehouse_sources.backend.facade.duckgres import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.facade.models import ExternalDataSchema from products.warehouse_sources.backend.facade.pipelines import DUCKGRES_BATCH_SINK_FLAG, is_duckgres_sink_team_member @@ -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=await database_sync_to_async(bind_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 56e43b68ed92..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 @@ -32,6 +32,7 @@ from products.managed_warehouse.backend.common import ( _get_org_id_for_team, duckgres_data_imports_schema, + duckgres_data_imports_table_name, get_config, get_duckgres_server_by_team_org, get_duckgres_server_for_organization, @@ -49,7 +50,6 @@ get_ducklake_register_data_imports_started_metric, record_ducklake_register_data_imports_stage_duration, ) -from products.warehouse_sources.backend.facade.duckgres import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.facade.models import ExternalDataSchema LOGGER = get_logger(__name__) @@ -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 = await database_sync_to_async(bind_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/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_common.py b/products/managed_warehouse/backend/tests/test_common.py index 4ebc6b93377a..fb230998e3b2 100644 --- a/products/managed_warehouse/backend/tests/test_common.py +++ b/products/managed_warehouse/backend/tests/test_common.py @@ -58,24 +58,14 @@ def test_region_follows_cloud_deployment(self, _name, deployment, expected): class TestDuckgresDataImportsTableName: @parameterized.expand( [ - ("copy_mysql", None, "copy_v1", "MySQL", "SalesEU", "orders", "mysql_saleseu_orders"), - ("copy_google_ads", None, "copy_v1", "GoogleAds", None, "video", "googleads_video"), - ("legacy_batch_tiktok", None, "legacy_batch_v1", "TikTokAds", None, "video", "tik_tok_ads_video"), - ( - "pinned_name", - "googleads_video_4f12abcd", - "legacy_batch_v1", - "GoogleAds", - None, - "video", - "googleads_video_4f12abcd", - ), + ("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_pinned_name_wins_and_null_uses_the_org_policy( + def test_uses_the_org_policy( self, _name: str, - pinned_name: str | None, naming_version: str, source_type: str, prefix: str | None, @@ -83,7 +73,6 @@ def test_pinned_name_wins_and_null_uses_the_org_policy( expected: str, ) -> None: schema = MagicMock() - schema.duckgres_table_name = pinned_name schema.source.source_type = source_type schema.source.prefix = prefix schema.normalized_name = normalized_name @@ -92,14 +81,9 @@ def test_pinned_name_wins_and_null_uses_the_org_policy( with patch( "products.managed_warehouse.backend.team_state.data_imports_table_naming_version", return_value=naming_version, - ) as mock_version: + ): assert duckgres_data_imports_table_name(schema) == expected - if pinned_name: - mock_version.assert_not_called() - else: - mock_version.assert_called_once_with(1) - TEST_CONFIG = { "DUCKLAKE_RDS_HOST": "localhost", diff --git a/products/managed_warehouse/backend/tests/test_cp_teams.py b/products/managed_warehouse/backend/tests/test_cp_teams.py index 7b0650bf01f1..8de7afb8ef25 100644 --- a/products/managed_warehouse/backend/tests/test_cp_teams.py +++ b/products/managed_warehouse/backend/tests/test_cp_teams.py @@ -214,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/warehouse_sources/backend/duckgres_naming.py b/products/warehouse_sources/backend/duckgres_naming.py index b66be6d3630b..cbbf904eb528 100644 --- a/products/warehouse_sources/backend/duckgres_naming.py +++ b/products/warehouse_sources/backend/duckgres_naming.py @@ -4,7 +4,7 @@ from products.warehouse_sources.backend.temporal.data_imports.naming_convention import NamingConvention -_IDENTIFIER_SANITIZE_RE = re.compile(r"[^A-Za-z0-9_]+") +_IDENTIFIER_SANITIZE_RE = re.compile(r"[^0-9a-zA-Z]+") _DUCKGRES_IDENTIFIER_MAX_LENGTH = 63 diff --git a/products/warehouse_sources/backend/duckgres_table_binding.py b/products/warehouse_sources/backend/duckgres_table_binding.py deleted file mode 100644 index 8a15877657de..000000000000 --- a/products/warehouse_sources/backend/duckgres_table_binding.py +++ /dev/null @@ -1,32 +0,0 @@ -from __future__ import annotations - -from django.db.models import Q - -from products.managed_warehouse.backend.facade import team_state -from products.warehouse_sources.backend.duckgres_naming import duckgres_data_imports_table_name_for_version -from products.warehouse_sources.backend.models.external_data_schema import ExternalDataSchema - - -def bind_duckgres_data_imports_table_name(schema: ExternalDataSchema) -> str: - pinned_name = schema.duckgres_table_name - if pinned_name: - return pinned_name - - naming_version = team_state.data_imports_table_naming_version(schema.team_id) - candidate = duckgres_data_imports_table_name_for_version( - schema.source.source_type, - schema.source.prefix, - schema.normalized_name, - naming_version, - ) - updated = ExternalDataSchema.objects.filter( - Q(duckgres_table_name__isnull=True) | Q(duckgres_table_name=""), pk=schema.pk - ).update(duckgres_table_name=candidate) - if updated: - schema.duckgres_table_name = candidate - return candidate - - schema.refresh_from_db(fields=["duckgres_table_name"]) - if not schema.duckgres_table_name: - raise RuntimeError(f"External data schema {schema.pk} has no Duckgres table name after binding") - return schema.duckgres_table_name diff --git a/products/warehouse_sources/backend/facade/duckgres.py b/products/warehouse_sources/backend/facade/duckgres.py index 33a5bc7c9b2d..91d4b4cfc485 100644 --- a/products/warehouse_sources/backend/facade/duckgres.py +++ b/products/warehouse_sources/backend/facade/duckgres.py @@ -1,15 +1,8 @@ from __future__ import annotations -from typing import TYPE_CHECKING - -if TYPE_CHECKING: - from products.warehouse_sources.backend.facade.models import ExternalDataSchema - - -def bind_duckgres_data_imports_table_name(schema: ExternalDataSchema) -> str: - from products.warehouse_sources.backend.duckgres_table_binding import bind_duckgres_data_imports_table_name - - return bind_duckgres_data_imports_table_name(schema) +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( @@ -18,6 +11,4 @@ def duckgres_data_imports_table_name_for_version( normalized_name: str, naming_version: str, ) -> str: - from products.warehouse_sources.backend.duckgres_naming import duckgres_data_imports_table_name_for_version - - return duckgres_data_imports_table_name_for_version(source_type, prefix, normalized_name, naming_version) + return _duckgres_data_imports_table_name_for_version(source_type, prefix, normalized_name, naming_version) diff --git a/products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py b/products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py deleted file mode 100644 index 4e00788abcb3..000000000000 --- a/products/warehouse_sources/backend/migrations/0116_add_duckgres_table_name.py +++ /dev/null @@ -1,15 +0,0 @@ -from django.db import migrations, models - - -class Migration(migrations.Migration): - dependencies = [ - ("warehouse_sources", "0115_scaffold_four_requested_sources"), - ] - - operations = [ - migrations.AddField( - model_name="externaldataschema", - name="duckgres_table_name", - field=models.CharField(blank=True, max_length=63, null=True), - ), - ] diff --git a/products/warehouse_sources/backend/models/external_data_schema.py b/products/warehouse_sources/backend/models/external_data_schema.py index 5cf93083df65..4729c67e6f9f 100644 --- a/products/warehouse_sources/backend/models/external_data_schema.py +++ b/products/warehouse_sources/backend/models/external_data_schema.py @@ -86,7 +86,6 @@ class SyncFrequency(models.TextChoices): # during multi-schema migration) to their original path. Empty for rows written before this # column existed — readers fall back to the legacy JSON key, then the normalized schema `name`. s3_folder_name = models.CharField(max_length=400, null=True, blank=True) - duckgres_table_name = models.CharField(max_length=63, null=True, blank=True) # Deprecated in favour of `sync_frequency_interval` sync_frequency = deprecate_field( models.CharField(max_length=128, choices=SyncFrequency, default=SyncFrequency.DAILY, blank=True) 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 0871380a16f8..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 @@ -26,7 +26,6 @@ get_duckgres_query_server_config, setup_duckgres_session, ) -from products.warehouse_sources.backend.duckgres_table_binding import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.models import ExternalDataJob, ExternalDataSchema from products.warehouse_sources.backend.temporal.data_imports.naming_convention import NamingConvention from products.warehouse_sources.backend.temporal.data_imports.pipelines.pipeline_v3.batch_consumer import ( @@ -246,8 +245,6 @@ def process_batch(batch: PendingBatch) -> None: raise ValueError(f"ExternalDataJob {batch.job_id} has no schema") schema = job.schema - bind_duckgres_data_imports_table_name(schema) - kind = "backfill" if _is_backfill_batch(batch) else "live" # One ORM lookup serves both the cache key and the connection config; it is # per-batch on purpose so the key always reflects the team's CURRENT org. 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 257bf7c53613..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 @@ -62,7 +67,6 @@ def _make_schema() -> Mock: schema = Mock() schema.team_id = 1 schema.normalized_name = "customers" - schema.duckgres_table_name = "stripe_customers" schema.source.source_type = "Stripe" schema.source.prefix = None return schema diff --git a/products/warehouse_sources/backend/tests/test_duckgres_naming.py b/products/warehouse_sources/backend/tests/test_duckgres_naming.py index d5c35d0614e0..8b016be3f48d 100644 --- a/products/warehouse_sources/backend/tests/test_duckgres_naming.py +++ b/products/warehouse_sources/backend/tests/test_duckgres_naming.py @@ -11,7 +11,7 @@ class TestDuckgresDataImportsTableNameForVersion: ("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", "ad_report", "tiktokads_prod_ad_report"), + ("tiktok_ads", "TikTokAds", "prod__us", "ad_report", "tiktokads_prod_us_ad_report"), ] ) def test_source_keys_are_stable_without_camel_case_splitting( diff --git a/products/warehouse_sources/backend/tests/test_models.py b/products/warehouse_sources/backend/tests/test_models.py index 951f2d3058c7..5cfc8f799136 100644 --- a/products/warehouse_sources/backend/tests/test_models.py +++ b/products/warehouse_sources/backend/tests/test_models.py @@ -16,7 +16,6 @@ from posthog.models.signals import model_activity_signal -from products.warehouse_sources.backend.duckgres_table_binding import bind_duckgres_data_imports_table_name from products.warehouse_sources.backend.models.credential import DataWarehouseCredential from products.warehouse_sources.backend.models.external_data_job import ExternalDataJob from products.warehouse_sources.backend.models.external_data_schema import ( @@ -77,20 +76,17 @@ def test_resolved_s3_folder_name( class TestExternalDataSchemaSave(BaseTest): - def _source(self, source_type: str = "Postgres", prefix: str | None = None) -> ExternalDataSource: + def _source(self) -> ExternalDataSource: return ExternalDataSource.objects.create( team_id=self.team.pk, source_id=str(uuid.uuid4()), connection_id=str(uuid.uuid4()), status="Completed", - source_type=source_type, - prefix=prefix, + source_type="Postgres", ) - def _create(self, name: str, *, source: ExternalDataSource | None = None, **kwargs) -> ExternalDataSchema: - return ExternalDataSchema.objects.create( - team_id=self.team.pk, source=source or self._source(), name=name, **kwargs - ) + def _create(self, name: str, **kwargs) -> ExternalDataSchema: + return ExternalDataSchema.objects.create(team_id=self.team.pk, source=self._source(), name=name, **kwargs) def test_save_populates_s3_folder_name_from_name(self) -> None: # The folder is the normalized name — never NULL for a new row. @@ -119,48 +115,6 @@ def test_partial_update_backfills_null_folder(self) -> None: schema.refresh_from_db() assert schema.s3_folder_name == "orders" - def test_new_schema_waits_for_a_duckgres_writer_to_pin_its_name(self) -> None: - schema = self._create("customer_orders", source=self._source("MySQL", "SalesEU")) - assert schema.duckgres_table_name is None - - @parameterized.expand( - [ - ("copy", "copy_v1", "tiktokads_ad_report"), - ("legacy_batch", "legacy_batch_v1", "tik_tok_ads_ad_report"), - ] - ) - def test_duckgres_writer_binds_the_org_naming_policy(self, _name: str, naming_version: str, expected: str) -> None: - schema = self._create("ad_report", source=self._source("TikTokAds")) - - with patch( - "products.warehouse_sources.backend.duckgres_table_binding.team_state.data_imports_table_naming_version", - return_value=naming_version, - ): - assert bind_duckgres_data_imports_table_name(schema) == expected - - schema.refresh_from_db() - assert schema.duckgres_table_name == expected - - def test_duckgres_writer_keeps_the_first_name_when_another_writer_wins(self) -> None: - schema = self._create("ad_report", source=self._source("TikTokAds")) - ExternalDataSchema.objects.filter(pk=schema.pk).update(duckgres_table_name="already_bound") - - with patch( - "products.warehouse_sources.backend.duckgres_table_binding.team_state.data_imports_table_naming_version", - return_value="copy_v1", - ): - assert bind_duckgres_data_imports_table_name(schema) == "already_bound" - - assert schema.duckgres_table_name == "already_bound" - - def test_existing_schema_without_pin_stays_on_legacy_naming(self) -> None: - schema = self._create("orders", source=self._source("MySQL", "SalesEU")) - schema.status = "Completed" - schema.save(update_fields=["status", "updated_at"]) - schema.refresh_from_db() - - assert schema.duckgres_table_name is None - class TestExternalDataSchemaActivityLogging(BaseTest): """Internal pipeline-driven bookkeeping saves must bypass ModelActivityMixin so they neither From 8e43fe8098c5ddea65e423cbf19914f230400166 Mon Sep 17 00:00:00 2001 From: eric Date: Sun, 2 Aug 2026 18:31:47 -0700 Subject: [PATCH 3/3] chore(data-warehouse): isolate client tests from control plane --- .../backend/tests/test_client.py | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) 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