diff --git a/README.md b/README.md index 8159aceb..fb1b5bb5 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 every data-import reader and writer derives the same physical table name. Operators can change the policy in the admin console or with `PUT /api/v1/orgs/:id` using `{"data_imports_table_naming_version":"copy_v1"}`. 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/admin/static/models.html b/controlplane/admin/static/models.html index bc5728d7..34b40c39 100644 --- a/controlplane/admin/static/models.html +++ b/controlplane/admin/static/models.html @@ -279,11 +279,11 @@ padding: 6px 18px; border-bottom: 1px solid var(--line); align-items: center; } .dform .frow label { color: var(--dim); font-size: 11px; letter-spacing: .04em; } - .dform .frow input { + .dform .frow input, .dform .frow select { background: var(--panel2); border: 1px solid var(--line-bright); border-radius: 4px; color: var(--text); font-family: inherit; font-size: 12px; padding: 5px 8px; width: 100%; } - .dform .frow input:focus { outline: none; border-color: var(--worker); } + .dform .frow input:focus, .dform .frow select:focus { outline: none; border-color: var(--worker); } .dform .factions { display: flex; gap: 8px; padding: 13px 18px; } .empty { padding: 40px 22px; color: var(--dim); font-size: 12px; letter-spacing: .04em; } @@ -725,7 +725,8 @@

Session Expired

// intersected with the live listing columns (orgEditableCols) so the edit form // auto-adapts as columns are added/removed from the model — non-mergeable // columns (name, database_name, users, warehouse, *_state, timestamps) are - // simply never in this allowlist. "int" → numeric input, else text. + // simply never in this allowlist. "int" uses a numeric input and + // "naming-version" uses the two supported policy values. const ORG_EDIT_FIELDS = { max_workers: "int", max_vcpus: "int", @@ -734,6 +735,7 @@

Session Expired

default_worker_ttl: "text", default_worker_min_hot_idle: "int", hostname_alias: "text", + data_imports_table_naming_version: "naming-version", }; function orgEditableCols() { const cols = (state.listing && state.listing.columns) || []; @@ -835,11 +837,18 @@

Session Expired

function orgEditFormHTML(row) { let html = '
'; for (const c of orgEditableCols()) { - const type = ORG_EDIT_FIELDS[c] === "int" ? "number" : "text"; const raw = row[c]; const v = (raw === null || raw === undefined) ? "" : String(raw); - html += '
" + - '
'; + if (ORG_EDIT_FIELDS[c] === "naming-version") { + const options = ["legacy_batch_v1", "copy_v1"].map(value => + '").join(""); + html += '
" + + '
"; + } else { + const type = ORG_EDIT_FIELDS[c] === "int" ? "number" : "text"; + html += '
" + + '
'; + } } html += '
' + '
'; diff --git a/controlplane/admin/ui/src/pages/OrgDetail.test.tsx b/controlplane/admin/ui/src/pages/OrgDetail.test.tsx index 0bc9e9c8..3f1a66dc 100644 --- a/controlplane/admin/ui/src/pages/OrgDetail.test.tsx +++ b/controlplane/admin/ui/src/pages/OrgDetail.test.tsx @@ -33,6 +33,7 @@ vi.mock("@/components/OrgTeamDialogs", () => ({ import { OrgDetail } from "./OrgDetail"; const warehouseUpdate = vi.fn(); +const orgUpdate = vi.fn(); const ok = (data: T) => ({ data, isSuccess: true, @@ -56,6 +57,7 @@ const ORG: Org = { default_worker_memory: "8Gi", default_worker_ttl: "75m", default_worker_min_hot_idle: 0, + data_imports_table_naming_version: "legacy_batch_v1", created_at: "2026-07-01T00:00:00Z", updated_at: "2026-07-01T00:00:00Z", }; @@ -116,7 +118,7 @@ function renderPage(metadataProxyEnabled: boolean) { ); } -describe("Org warehouse metadata proxy setting", () => { +describe("Org detail", () => { beforeEach(() => { vi.clearAllMocks(); identity.useIdentity.mockReturnValue({ @@ -124,7 +126,7 @@ describe("Org warehouse metadata proxy setting", () => { me: { email: "admin@example.com", role: "admin", source: "sso" }, }); hooks.useOrg.mockReturnValue(ok(ORG)); - hooks.useUpdateOrg.mockReturnValue(mut()); + hooks.useUpdateOrg.mockReturnValue(mut(orgUpdate)); hooks.useDeleteOrg.mockReturnValue(mut()); hooks.useUpdateWarehouse.mockReturnValue(mut(warehouseUpdate)); hooks.useDeprovisionWarehouse.mockReturnValue(mut()); @@ -150,4 +152,22 @@ describe("Org warehouse metadata proxy setting", () => { expect(warehouseUpdate).toHaveBeenCalledTimes(1); expect(warehouseUpdate).toHaveBeenCalledWith({ metadata_proxy_enabled: next }); }); + + it("saves a changed data import table naming version", async () => { + const user = userEvent.setup(); + HTMLElement.prototype.scrollIntoView = vi.fn(); + renderPage(false); + + const namingSelect = screen.getByLabelText("Data import table naming"); + expect(namingSelect).toHaveTextContent("Legacy batch (legacy_batch_v1)"); + + namingSelect.focus(); + await user.keyboard("{Enter}{ArrowDown}{Enter}"); + await user.click(screen.getByText("Save changes")); + + expect(orgUpdate).toHaveBeenCalledTimes(1); + expect(orgUpdate).toHaveBeenCalledWith( + expect.objectContaining({ data_imports_table_naming_version: "copy_v1" }), + ); + }); }); diff --git a/controlplane/admin/ui/src/pages/OrgDetail.tsx b/controlplane/admin/ui/src/pages/OrgDetail.tsx index 69cb56bc..cd9ef43e 100644 --- a/controlplane/admin/ui/src/pages/OrgDetail.tsx +++ b/controlplane/admin/ui/src/pages/OrgDetail.tsx @@ -8,6 +8,7 @@ import { Input } from "@/components/ui/input"; import { Label } from "@/components/ui/label"; import { Badge } from "@/components/ui/badge"; import { Switch } from "@/components/ui/switch"; +import { Select, SelectContent, SelectItem, SelectTrigger, SelectValue } from "@/components/ui/select"; import { StateBadge } from "@/components/StateBadge"; import { Tooltip, TooltipContent, TooltipTrigger } from "@/components/ui/tooltip"; import { AdminGate } from "@/components/AdminOnly"; @@ -44,7 +45,7 @@ import { LegacyNamesBadge, } from "@/components/OrgTeamDialogs"; import { Table, TableBody, TableCell, TableHead, TableHeader, TableRow } from "@/components/ui/table"; -import type { ManagedWarehouse, OrgTeam, OrgUpdate } from "@/types/api"; +import type { DataImportsTableNamingVersion, ManagedWarehouse, OrgTeam, OrgUpdate } from "@/types/api"; interface FormState { max_workers: string; @@ -54,6 +55,7 @@ interface FormState { default_worker_ttl: string; default_worker_min_hot_idle: string; hostname_alias: string; + data_imports_table_naming_version: DataImportsTableNamingVersion; } function orgToForm(o: { @@ -64,6 +66,7 @@ function orgToForm(o: { default_worker_ttl: string; default_worker_min_hot_idle: number; hostname_alias: string | null; + data_imports_table_naming_version: DataImportsTableNamingVersion; }): FormState { return { max_workers: String(o.max_workers), @@ -73,6 +76,7 @@ function orgToForm(o: { default_worker_ttl: o.default_worker_ttl, default_worker_min_hot_idle: String(o.default_worker_min_hot_idle), hostname_alias: o.hostname_alias ?? "", + data_imports_table_naming_version: o.data_imports_table_naming_version, }; } @@ -124,6 +128,7 @@ export function OrgDetail() { default_worker_ttl: form.default_worker_ttl, default_worker_min_hot_idle: Number(form.default_worker_min_hot_idle) || 0, hostname_alias: form.hostname_alias === "" ? "" : form.hostname_alias, + data_imports_table_naming_version: form.data_imports_table_naming_version, }; try { await update.mutateAsync(body); @@ -244,6 +249,35 @@ export function OrgDetail() { onChange={(e) => set("hostname_alias", e.target.value)} /> + + + + {form.data_imports_table_naming_version !== org.data?.data_imports_table_naming_version && ( +

+ + + Changing the naming version does not rename or move existing data. Migrate the + existing tables before saving this change. + +

+ )}