From f92f3220680ae66f4062a6af067145ae84d5a448 Mon Sep 17 00:00:00 2001 From: eric Date: Sat, 1 Aug 2026 21:47:42 -0700 Subject: [PATCH 1/2] feat(control-plane): version data import table naming --- README.md | 1 + .../configstore/data_imports_table_naming.go | 6 ++ .../000034_org_data_imports_table_naming.sql | 17 ++++ controlplane/configstore/models.go | 17 ++-- controlplane/provisioning/api.go | 14 ++- controlplane/provisioning/api_test.go | 11 ++- tests/configstore/migrations_postgres_test.go | 88 ++++++++++++++++++- 7 files changed, 139 insertions(+), 15 deletions(-) create mode 100644 controlplane/configstore/data_imports_table_naming.go create mode 100644 controlplane/configstore/migrations/000034_org_data_imports_table_naming.sql diff --git a/README.md b/README.md index 8159aceb..e003df66 100644 --- a/README.md +++ b/README.md @@ -952,6 +952,7 @@ The shared K8s pool spawns workers on-demand, reserves them per org, activates t Managed-warehouse contract notes: - At most one managed-warehouse row exists per team. The row may be absent before first provisioning or after cleanup, but there is never more than one active warehouse contract for a team. +- Each org has a `data_imports_table_naming_version`. Migration `000034` assigns `legacy_batch_v1` to orgs that already exist and changes the database default to `copy_v1` for orgs created afterward. `GET /api/v1/orgs/:id/teams` returns the org-level value alongside the team rows so data-import writers can pin one physical table name consistently. - The admin API exposes that contract at `GET /api/v1/teams/:name/warehouse` and `PUT /api/v1/teams/:name/warehouse`. Team list/get responses also include a nested `warehouse` object when present. - Org rows support optional `max_vcpus` on `POST /api/v1/orgs` and `PUT /api/v1/orgs/:id`. In K8s multi-tenant mode, this caps the org's active admitted worker pod vCPUs; `0` means unlimited. - User rows support an optional `max_vcpus` field on `POST /api/v1/users` and `PUT /api/v1/orgs/:id/users/:username`. `max_vcpus` limits the user's active admitted worker pod vCPUs in K8s multi-tenant mode; `0` means unlimited. diff --git a/controlplane/configstore/data_imports_table_naming.go b/controlplane/configstore/data_imports_table_naming.go new file mode 100644 index 00000000..ce835c5d --- /dev/null +++ b/controlplane/configstore/data_imports_table_naming.go @@ -0,0 +1,6 @@ +package configstore + +const ( + DataImportsTableNamingVersionLegacyBatchV1 = "legacy_batch_v1" + DataImportsTableNamingVersionCopyV1 = "copy_v1" +) diff --git a/controlplane/configstore/migrations/000034_org_data_imports_table_naming.sql b/controlplane/configstore/migrations/000034_org_data_imports_table_naming.sql new file mode 100644 index 00000000..311ea58e --- /dev/null +++ b/controlplane/configstore/migrations/000034_org_data_imports_table_naming.sql @@ -0,0 +1,17 @@ +-- +goose Up +ALTER TABLE duckgres_orgs + ADD COLUMN IF NOT EXISTS data_imports_table_naming_version VARCHAR(32) NOT NULL DEFAULT 'legacy_batch_v1'; + +ALTER TABLE duckgres_orgs + ALTER COLUMN data_imports_table_naming_version SET DEFAULT 'copy_v1'; + +ALTER TABLE duckgres_orgs + ADD CONSTRAINT duckgres_orgs_data_imports_table_naming_version_check + CHECK (data_imports_table_naming_version IN ('legacy_batch_v1', 'copy_v1')); + +-- +goose Down +ALTER TABLE duckgres_orgs + DROP CONSTRAINT IF EXISTS duckgres_orgs_data_imports_table_naming_version_check; + +ALTER TABLE duckgres_orgs + DROP COLUMN IF EXISTS data_imports_table_naming_version; diff --git a/controlplane/configstore/models.go b/controlplane/configstore/models.go index d8ae47ba..cbce38b8 100644 --- a/controlplane/configstore/models.go +++ b/controlplane/configstore/models.go @@ -21,14 +21,15 @@ type Org struct { // human editability) applied to connections that don't size themselves via // the duckgres.worker_* startup options. Empty = unset. Versioned SQL // migrations add these columns. - DefaultWorkerCPU string `gorm:"size:32" json:"default_worker_cpu"` - DefaultWorkerMemory string `gorm:"size:32" json:"default_worker_memory"` - DefaultWorkerTTL string `gorm:"size:32" json:"default_worker_ttl"` - DefaultWorkerMinHotIdle int `gorm:"default:0" json:"default_worker_min_hot_idle"` - Teams []OrgTeam `gorm:"foreignKey:OrgID;references:Name;constraint:OnDelete:CASCADE" json:"teams,omitempty"` - Users []OrgUser `gorm:"foreignKey:OrgID;references:Name" json:"users,omitempty"` - Warehouse *ManagedWarehouse `gorm:"foreignKey:OrgID;references:Name;constraint:OnDelete:CASCADE" json:"warehouse,omitempty"` - CreatedAt time.Time `json:"created_at"` + DefaultWorkerCPU string `gorm:"size:32" json:"default_worker_cpu"` + DefaultWorkerMemory string `gorm:"size:32" json:"default_worker_memory"` + DefaultWorkerTTL string `gorm:"size:32" json:"default_worker_ttl"` + DefaultWorkerMinHotIdle int `gorm:"default:0" json:"default_worker_min_hot_idle"` + DataImportsTableNamingVersion string `gorm:"size:32;not null;default:copy_v1" json:"data_imports_table_naming_version"` + Teams []OrgTeam `gorm:"foreignKey:OrgID;references:Name;constraint:OnDelete:CASCADE" json:"teams,omitempty"` + Users []OrgUser `gorm:"foreignKey:OrgID;references:Name" json:"users,omitempty"` + Warehouse *ManagedWarehouse `gorm:"foreignKey:OrgID;references:Name;constraint:OnDelete:CASCADE" json:"warehouse,omitempty"` + CreatedAt time.Time `json:"created_at"` // UpdatedAt doubles as an input to the discovery change marker // (ConfigStore.LatestConfigChange): DeleteOrgTeamTx touches it so a // team-row DELETE — which leaves no updated_at of its own behind — diff --git a/controlplane/provisioning/api.go b/controlplane/provisioning/api.go index d1449c14..4695ddc0 100644 --- a/controlplane/provisioning/api.go +++ b/controlplane/provisioning/api.go @@ -602,6 +602,15 @@ type orgTeamUpsertRequest struct { func (h *handler) listOrgTeams(c *gin.Context) { orgID := c.Param("id") + org, err := h.store.GetOrg(orgID) + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + c.JSON(http.StatusNotFound, gin.H{"error": "org not found"}) + return + } + c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) + return + } teams, err := h.store.ListOrgTeams(orgID) if err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { @@ -611,7 +620,10 @@ func (h *handler) listOrgTeams(c *gin.Context) { c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } - c.JSON(http.StatusOK, gin.H{"teams": teams}) + c.JSON(http.StatusOK, gin.H{ + "teams": teams, + "data_imports_table_naming_version": org.DataImportsTableNamingVersion, + }) } // upsertOrgTeam creates or overwrites one (org, team) row. This endpoint IS diff --git a/controlplane/provisioning/api_test.go b/controlplane/provisioning/api_test.go index c5b066af..45b498f3 100644 --- a/controlplane/provisioning/api_test.go +++ b/controlplane/provisioning/api_test.go @@ -1174,7 +1174,10 @@ func TestOrgTeamUpsertValidation(t *testing.T) { func TestOrgTeamUpsertCreatesAndLists(t *testing.T) { store := newFakeStore() - store.orgs["acme"] = &configstore.Org{Name: "acme"} + store.orgs["acme"] = &configstore.Org{ + Name: "acme", + DataImportsTableNamingVersion: configstore.DataImportsTableNamingVersionCopyV1, + } router := newTestRouter(store) rec := doJSON(t, router, http.MethodPost, "/api/v1/orgs/acme/teams", @@ -1198,7 +1201,8 @@ func TestOrgTeamUpsertCreatesAndLists(t *testing.T) { t.Fatalf("list status = %d, want 200: %s", rec.Code, rec.Body.String()) } var listing struct { - Teams []configstore.OrgTeam `json:"teams"` + Teams []configstore.OrgTeam `json:"teams"` + DataImportsTableNamingVersion string `json:"data_imports_table_naming_version"` } if err := json.Unmarshal(rec.Body.Bytes(), &listing); err != nil { t.Fatalf("decode listing: %v", err) @@ -1206,6 +1210,9 @@ func TestOrgTeamUpsertCreatesAndLists(t *testing.T) { if len(listing.Teams) != 1 || listing.Teams[0].TeamID != 7 { t.Fatalf("listing = %+v, want the created team", listing.Teams) } + if listing.DataImportsTableNamingVersion != configstore.DataImportsTableNamingVersionCopyV1 { + t.Fatalf("data imports naming version = %q, want copy_v1", listing.DataImportsTableNamingVersion) + } } func TestOrgTeamListUnknownOrg404(t *testing.T) { diff --git a/tests/configstore/migrations_postgres_test.go b/tests/configstore/migrations_postgres_test.go index 813845f0..87c6b39b 100644 --- a/tests/configstore/migrations_postgres_test.go +++ b/tests/configstore/migrations_postgres_test.go @@ -52,7 +52,8 @@ func TestConfigStoreRunsVersionedSQLMigrations(t *testing.T) { requireGooseMigrationRecorded(t, db, 31) requireGooseMigrationRecorded(t, db, 32) requireGooseMigrationRecorded(t, db, 33) - requireGooseLatestVersion(t, db, 33) + requireGooseMigrationRecorded(t, db, 34) + requireGooseLatestVersion(t, db, 34) requireTableAbsent(t, db, "duckgres_schema_migrations") // Migration 000018 added the reshard operation + verbose log tables. @@ -120,6 +121,11 @@ func TestConfigStoreRunsVersionedSQLMigrations(t *testing.T) { requireColumnNullable(t, db, "duckgres_org_teams", "schema_data_imports_name") requireUniqueIndex(t, db, "duckgres_org_teams", "org_id,schema_name") + // Migration 000034 preserves the table naming used by existing orgs while + // selecting the copy workflow naming for orgs created after deployment. + requireColumnNotNull(t, db, "duckgres_orgs", "data_imports_table_naming_version") + requireColumnDefault(t, db, "duckgres_orgs", "data_imports_table_naming_version", "'copy_v1'::character varying") + // Migration 000026 added PostHog's cached earliest-event date (nullable // DATE — NULL until the PostHog sensor resolves it). requireColumnNullable(t, db, "duckgres_org_teams", "earliest_event_date") @@ -207,8 +213,8 @@ func TestConfigStoreSQLMigrationsUpgradeVersion8Schema(t *testing.T) { t.Cleanup(func() { _ = baselineDB.Close() }) - if err := store.DB().Exec(` + ALTER TABLE duckgres_orgs DROP COLUMN data_imports_table_naming_version; ALTER TABLE duckgres_orgs DROP COLUMN max_vcpus; ALTER TABLE duckgres_org_users DROP COLUMN max_vcpus; ALTER TABLE duckgres_org_users DROP COLUMN disabled; @@ -237,7 +243,7 @@ func TestConfigStoreSQLMigrationsUpgradeVersion8Schema(t *testing.T) { ); DROP TABLE IF EXISTS duckgres_reshard_operation_log; DROP TABLE IF EXISTS duckgres_reshard_operations; - DELETE FROM goose_db_version WHERE version_id IN (9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20, 21, 22, 23, 24, 25, 26, 27, 28, 29, 30, 31, 32, 33); + DELETE FROM goose_db_version WHERE version_id IN (9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20, 21, 22, 23, 24, 25, 26, 27, 28, 29, 30, 31, 32, 33, 34); `).Error; err != nil { t.Fatalf("downgrade baseline schema to pre-v9 shape: %v", err) } @@ -284,7 +290,8 @@ func TestConfigStoreSQLMigrationsUpgradeVersion8Schema(t *testing.T) { requireGooseMigrationRecorded(t, upgradedDB, 31) requireGooseMigrationRecorded(t, upgradedDB, 32) requireGooseMigrationRecorded(t, upgradedDB, 33) - requireGooseLatestVersion(t, upgradedDB, 33) + requireGooseMigrationRecorded(t, upgradedDB, 34) + requireGooseLatestVersion(t, upgradedDB, 34) requireColumnPresent(t, upgradedDB, "duckgres_reshard_operations", "password_url") requireTablePresent(t, upgradedDB, "duckgres_worker_spawn_log") requireColumnDefault(t, upgradedDB, "duckgres_orgs", "max_vcpus", "0") @@ -300,6 +307,68 @@ func TestConfigStoreSQLMigrationsUpgradeVersion8Schema(t *testing.T) { requireColumnAbsent(t, upgradedDB, "duckgres_managed_warehouses", "iceberg_enabled") requireColumnDefault(t, upgradedDB, "duckgres_managed_warehouses", "metadata_proxy_enabled", "false") requireColumnAbsent(t, upgradedDB, "duckgres_org_users", "default_catalog") + +} + +func TestConfigStoreSQLMigration34VersionsExistingAndNewOrgs(t *testing.T) { + _, connStr := newIsolatedConfigStoreSchema(t) + store, err := cpconfigStoreNew(connStr) + if err != nil { + t.Fatalf("create baseline config store: %v", err) + } + baselineDB := storeDB(t, store) + t.Cleanup(func() { + _ = baselineDB.Close() + }) + + if err := store.DB().Exec(` + INSERT INTO duckgres_orgs (name, database_name, created_at, updated_at) + VALUES ('existing-naming-policy', 'existing-naming-policy', now(), now()); + ALTER TABLE duckgres_orgs DROP COLUMN data_imports_table_naming_version; + DELETE FROM goose_db_version WHERE version_id = 34; + `).Error; err != nil { + t.Fatalf("restore pre-migration-34 schema: %v", err) + } + requireGooseLatestVersion(t, baselineDB, 33) + + upgradedStore, err := cpconfigStoreNew(connStr) + if err != nil { + t.Fatalf("apply migration 34: %v", err) + } + upgradedDB := storeDB(t, upgradedStore) + t.Cleanup(func() { + _ = upgradedDB.Close() + }) + + var existingNamingVersion string + if err := upgradedStore.DB().Raw(` + SELECT data_imports_table_naming_version + FROM duckgres_orgs + WHERE name = 'existing-naming-policy' + `).Scan(&existingNamingVersion).Error; err != nil { + t.Fatalf("read existing org naming version: %v", err) + } + if existingNamingVersion != cpconfigstore.DataImportsTableNamingVersionLegacyBatchV1 { + t.Fatalf("existing org naming version = %q, want legacy_batch_v1", existingNamingVersion) + } + + if err := upgradedStore.DB().Exec(` + INSERT INTO duckgres_orgs (name, database_name, created_at, updated_at) + VALUES ('new-naming-policy', 'new-naming-policy', now(), now()) + `).Error; err != nil { + t.Fatalf("create org after naming migration: %v", err) + } + var newNamingVersion string + if err := upgradedStore.DB().Raw(` + SELECT data_imports_table_naming_version + FROM duckgres_orgs + WHERE name = 'new-naming-policy' + `).Scan(&newNamingVersion).Error; err != nil { + t.Fatalf("read new org naming version: %v", err) + } + if newNamingVersion != cpconfigstore.DataImportsTableNamingVersionCopyV1 { + t.Fatalf("new org naming version = %q, want copy_v1", newNamingVersion) + } } func TestConfigStoreSQLMigrationsUpgradeOldOrgSchema(t *testing.T) { @@ -368,6 +437,17 @@ func TestConfigStoreSQLMigrationsUpgradeOldOrgSchema(t *testing.T) { } requireColumnAbsent(t, sqlDB, "duckgres_orgs", "max_connections") requireGooseMigrationRecorded(t, sqlDB, 3) + var namingVersion string + if err := store.DB().Raw(` + SELECT data_imports_table_naming_version + FROM duckgres_orgs + WHERE name = 'old-org' + `).Scan(&namingVersion).Error; err != nil { + t.Fatalf("read migrated data imports naming version: %v", err) + } + if namingVersion != cpconfigstore.DataImportsTableNamingVersionLegacyBatchV1 { + t.Fatalf("migrated data imports naming version = %q, want legacy_batch_v1", namingVersion) + } // Migration 000024 backfilled the legacy default_team_id value into the // org's team row and dropped the column. From e4a4f546d75f0e0b2279932a47709e72c9d6d356 Mon Sep 17 00:00:00 2001 From: eric Date: Sun, 2 Aug 2026 18:08:08 -0700 Subject: [PATCH 2/2] feat(control-plane): make naming policy configurable --- README.md | 2 +- controlplane/admin/api.go | 29 +++++ controlplane/admin/api_postgres_test.go | 25 +++++ controlplane/admin/api_test.go | 102 ++++++++++++++++++ .../configstore/data_imports_table_naming.go | 9 ++ 5 files changed, 166 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index e003df66..a4d17478 100644 --- a/README.md +++ b/README.md @@ -952,7 +952,7 @@ The shared K8s pool spawns workers on-demand, reserves them per org, activates t Managed-warehouse contract notes: - At most one managed-warehouse row exists per team. The row may be absent before first provisioning or after cleanup, but there is never more than one active warehouse contract for a team. -- Each org has a `data_imports_table_naming_version`. Migration `000034` assigns `legacy_batch_v1` to orgs that already exist and changes the database default to `copy_v1` for orgs created afterward. `GET /api/v1/orgs/:id/teams` returns the org-level value alongside the team rows so data-import writers can pin one physical table name consistently. +- Each org has a `data_imports_table_naming_version`. Migration `000034` assigns `legacy_batch_v1` to orgs that already exist and changes the database default to `copy_v1` for orgs created afterward. `GET /api/v1/orgs/:id/teams` returns the org-level value alongside the team rows so every data-import reader and writer derives the same physical table name. Operators can change the policy with `PUT /api/v1/orgs/:id`; migrate existing tables before changing an org that has already written data. - The admin API exposes that contract at `GET /api/v1/teams/:name/warehouse` and `PUT /api/v1/teams/:name/warehouse`. Team list/get responses also include a nested `warehouse` object when present. - Org rows support optional `max_vcpus` on `POST /api/v1/orgs` and `PUT /api/v1/orgs/:id`. In K8s multi-tenant mode, this caps the org's active admitted worker pod vCPUs; `0` means unlimited. - User rows support an optional `max_vcpus` field on `POST /api/v1/users` and `PUT /api/v1/orgs/:id/users/:username`. `max_vcpus` limits the user's active admitted worker pod vCPUs in K8s multi-tenant mode; `0` means unlimited. diff --git a/controlplane/admin/api.go b/controlplane/admin/api.go index 2220fbb5..3b7fd763 100644 --- a/controlplane/admin/api.go +++ b/controlplane/admin/api.go @@ -262,6 +262,9 @@ func (s *gormAPIStore) UpdateOrg(name string, updates configstore.Org) (*configs "default_worker_ttl": updates.DefaultWorkerTTL, "default_worker_min_hot_idle": updates.DefaultWorkerMinHotIdle, } + if updates.DataImportsTableNamingVersion != "" { + fields["data_imports_table_naming_version"] = updates.DataImportsTableNamingVersion + } // HostnameAlias is *string: nil = preserve, "" = clear (NULL), "x" = set. // NULL releases the unique-index slot so other orgs can take that alias. if updates.HostnameAlias != nil { @@ -937,6 +940,12 @@ func (h *apiHandler) updateOrg(c *gin.Context) { c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) return } + if _, ok := fields["data_imports_table_naming_version"]; ok { + if err := validateDataImportsTableNamingVersion(updates.DataImportsTableNamingVersion); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + } existing, err := h.store.GetOrg(name) if err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { @@ -970,6 +979,9 @@ func (h *apiHandler) updateOrg(c *gin.Context) { if _, ok := fields["hostname_alias"]; ok { merged.HostnameAlias = updates.HostnameAlias } + if _, ok := fields["data_imports_table_naming_version"]; ok { + merged.DataImportsTableNamingVersion = updates.DataImportsTableNamingVersion + } // Audit detail: which fields changed and their old → new values, so the // console shows "max_workers 4 → 10" instead of a bare "org.update". These // are all non-sensitive config columns (no credentials among them). @@ -986,6 +998,7 @@ func (h *apiHandler) updateOrg(c *gin.Context) { addChange("default_worker_ttl", orgStr(existing.DefaultWorkerTTL), orgStr(merged.DefaultWorkerTTL)) addChange("default_worker_min_hot_idle", existing.DefaultWorkerMinHotIdle, merged.DefaultWorkerMinHotIdle) addChange("hostname_alias", orgStrPtr(existing.HostnameAlias), orgStrPtr(merged.HostnameAlias)) + addChange("data_imports_table_naming_version", existing.DataImportsTableNamingVersion, merged.DataImportsTableNamingVersion) if len(changes) > 0 { setAuditDetail(c, strings.Join(changes, ", ")) } @@ -1307,6 +1320,11 @@ func validateOrgMutationPayload(org *configstore.Org) error { return err } } + if org.DataImportsTableNamingVersion != "" { + if err := validateDataImportsTableNamingVersion(org.DataImportsTableNamingVersion); err != nil { + return err + } + } if err := validateOrgDefaultWorkerProfile(org); err != nil { return err } @@ -1319,6 +1337,17 @@ func validateOrgMutationPayload(org *configstore.Org) error { return nil } +func validateDataImportsTableNamingVersion(version string) error { + if configstore.IsValidDataImportsTableNamingVersion(version) { + return nil + } + return fmt.Errorf( + "data_imports_table_naming_version must be %q or %q", + configstore.DataImportsTableNamingVersionLegacyBatchV1, + configstore.DataImportsTableNamingVersionCopyV1, + ) +} + // validateOrgDefaultWorkerProfile rejects garbage default-worker-profile // values at the API boundary so they can never enter the config store (the // control plane tolerates bad rows by ignoring them, but a 400 here surfaces diff --git a/controlplane/admin/api_postgres_test.go b/controlplane/admin/api_postgres_test.go index a2c2ffeb..4fb1e551 100644 --- a/controlplane/admin/api_postgres_test.go +++ b/controlplane/admin/api_postgres_test.go @@ -104,6 +104,31 @@ func TestAdminAdmissionConfigMutationsSerializePostgres(t *testing.T) { } } +func TestAdminUpdateOrgPersistsDataImportsTableNamingVersionPostgres(t *testing.T) { + store := newPostgresConfigStore(t) + if err := store.DB().Create(&configstore.Org{ + Name: "naming-policy-org", + DatabaseName: "naming_policy_org", + DataImportsTableNamingVersion: configstore.DataImportsTableNamingVersionLegacyBatchV1, + }).Error; err != nil { + t.Fatalf("create org: %v", err) + } + + apiStore := newGormAPIStore(store).(*gormAPIStore) + updated, found, err := apiStore.UpdateOrg("naming-policy-org", configstore.Org{ + DataImportsTableNamingVersion: configstore.DataImportsTableNamingVersionCopyV1, + }) + if err != nil { + t.Fatalf("update org: %v", err) + } + if !found { + t.Fatal("updated org was not found") + } + if updated.DataImportsTableNamingVersion != configstore.DataImportsTableNamingVersionCopyV1 { + t.Fatalf("stored naming version = %q, want copy_v1", updated.DataImportsTableNamingVersion) + } +} + func seedAdmissionMutationUser(t *testing.T, store *configstore.ConfigStore, orgID, username string) { t.Helper() if err := store.DB().Create(&configstore.Org{Name: orgID, DatabaseName: orgID}).Error; err != nil { diff --git a/controlplane/admin/api_test.go b/controlplane/admin/api_test.go index 19d669ee..531399a6 100644 --- a/controlplane/admin/api_test.go +++ b/controlplane/admin/api_test.go @@ -83,6 +83,9 @@ func (s *fakeAPIStore) UpdateOrg(name string, updates configstore.Org) (*configs org.DefaultWorkerMemory = updates.DefaultWorkerMemory org.DefaultWorkerTTL = updates.DefaultWorkerTTL org.DefaultWorkerMinHotIdle = updates.DefaultWorkerMinHotIdle + if updates.DataImportsTableNamingVersion != "" { + org.DataImportsTableNamingVersion = updates.DataImportsTableNamingVersion + } if updates.HostnameAlias != nil { if *updates.HostnameAlias == "" { org.HostnameAlias = nil @@ -2349,6 +2352,105 @@ func TestUpdateOrgRejectsNegativeMaxVCPUs(t *testing.T) { } } +func TestUpdateOrgDataImportsTableNamingVersion(t *testing.T) { + for _, tc := range []struct { + name string + from string + to string + }{ + {"legacy_to_copy", configstore.DataImportsTableNamingVersionLegacyBatchV1, configstore.DataImportsTableNamingVersionCopyV1}, + {"copy_to_legacy", configstore.DataImportsTableNamingVersionCopyV1, configstore.DataImportsTableNamingVersionLegacyBatchV1}, + } { + t.Run(tc.name, func(t *testing.T) { + store := newFakeAPIStore() + store.orgs["analytics"] = &configstore.Org{ + Name: "analytics", + DataImportsTableNamingVersion: tc.from, + } + + gin.SetMode(gin.TestMode) + router := gin.New() + var detail string + router.Use(func(c *gin.Context) { + c.Next() + detail = c.GetString(ctxAuditDetailKey) + }) + registerAPIWithStore(router.Group("/api/v1"), store, nil, nil) + + rec := adminJSON(t, router, http.MethodPut, "/api/v1/orgs/analytics", + fmt.Sprintf(`{"data_imports_table_naming_version":%q}`, tc.to)) + if rec.Code != http.StatusOK { + t.Fatalf("status = %d, want %d: %s", rec.Code, http.StatusOK, rec.Body.String()) + } + var response configstore.Org + if err := json.Unmarshal(rec.Body.Bytes(), &response); err != nil { + t.Fatalf("decode response: %v", err) + } + if response.DataImportsTableNamingVersion != tc.to { + t.Fatalf("response naming version = %q, want %q", response.DataImportsTableNamingVersion, tc.to) + } + if got := store.orgs["analytics"].DataImportsTableNamingVersion; got != tc.to { + t.Fatalf("stored naming version = %q, want %q", got, tc.to) + } + wantDetail := fmt.Sprintf("data_imports_table_naming_version %s → %s", tc.from, tc.to) + if !strings.Contains(detail, wantDetail) { + t.Fatalf("audit detail = %q, want %q", detail, wantDetail) + } + + rec = adminJSON(t, router, http.MethodPut, "/api/v1/orgs/analytics", `{"max_workers":3}`) + if rec.Code != http.StatusOK { + t.Fatalf("unrelated update status = %d, want %d: %s", rec.Code, http.StatusOK, rec.Body.String()) + } + if got := store.orgs["analytics"].DataImportsTableNamingVersion; got != tc.to { + t.Fatalf("unrelated update changed naming version to %q", got) + } + }) + } +} + +func TestUpdateOrgRejectsInvalidDataImportsTableNamingVersion(t *testing.T) { + for _, tc := range []struct { + name string + value string + }{ + {"empty", `""`}, + {"null", `null`}, + {"unknown", `"future_v1"`}, + } { + t.Run(tc.name, func(t *testing.T) { + store := newFakeAPIStore() + store.orgs["analytics"] = &configstore.Org{ + Name: "analytics", + DataImportsTableNamingVersion: configstore.DataImportsTableNamingVersionLegacyBatchV1, + } + router := newTestAPIRouter(store) + + rec := adminJSON(t, router, http.MethodPut, "/api/v1/orgs/analytics", + fmt.Sprintf(`{"data_imports_table_naming_version":%s}`, tc.value)) + if rec.Code != http.StatusBadRequest { + t.Fatalf("status = %d, want %d: %s", rec.Code, http.StatusBadRequest, rec.Body.String()) + } + if got := store.orgs["analytics"].DataImportsTableNamingVersion; got != configstore.DataImportsTableNamingVersionLegacyBatchV1 { + t.Fatalf("invalid update changed naming version to %q", got) + } + }) + } +} + +func TestCreateOrgRejectsInvalidDataImportsTableNamingVersion(t *testing.T) { + store := newFakeAPIStore() + router := newTestAPIRouter(store) + + rec := adminJSON(t, router, http.MethodPost, "/api/v1/orgs", + `{"name":"analytics","database_name":"analytics","team_id":1,"data_imports_table_naming_version":"future_v1"}`) + if rec.Code != http.StatusBadRequest { + t.Fatalf("status = %d, want %d: %s", rec.Code, http.StatusBadRequest, rec.Body.String()) + } + if _, ok := store.orgs["analytics"]; ok { + t.Fatal("org must not be created with an invalid naming version") + } +} + // --- Org default worker profile (default_worker_cpu/memory/ttl) --- func TestUpdateOrgSetsDefaultWorkerProfile(t *testing.T) { diff --git a/controlplane/configstore/data_imports_table_naming.go b/controlplane/configstore/data_imports_table_naming.go index ce835c5d..32b6b4bd 100644 --- a/controlplane/configstore/data_imports_table_naming.go +++ b/controlplane/configstore/data_imports_table_naming.go @@ -4,3 +4,12 @@ const ( DataImportsTableNamingVersionLegacyBatchV1 = "legacy_batch_v1" DataImportsTableNamingVersionCopyV1 = "copy_v1" ) + +func IsValidDataImportsTableNamingVersion(version string) bool { + switch version { + case DataImportsTableNamingVersionLegacyBatchV1, DataImportsTableNamingVersionCopyV1: + return true + default: + return false + } +}