From 23c3e598be775ecdc86fad1e4d029517bf8e79b0 Mon Sep 17 00:00:00 2001 From: rajdangwal Date: Wed, 22 Jul 2026 15:43:40 +0530 Subject: [PATCH 1/4] USTORAGE-27286: Add metadata-column-naming-scheme flag to tableflow topic commands Add a --metadata-column-naming-scheme flag to `tableflow topic enable` (create) and `tableflow topic update`, sending it on the config only when the user provides it, and surface the value in describe/list/enable/update output (omitted when unset). Bumps ccloud-sdk-go-v2/tableflow to v0.7.0, which adds the field to the topic config spec. Adds integration test cases for enable/update with the flag (with a mock update handler that echoes the scheme) and regenerates the affected help goldens. --- go.mod | 4 +- go.sum | 8 +- internal/tableflow/command_topic.go | 84 ++++++++++--------- internal/tableflow/command_topic_enable.go | 10 +++ internal/tableflow/command_topic_list.go | 35 ++++---- internal/tableflow/command_topic_update.go | 10 +++ .../output/tableflow/topic/enable-help.golden | 29 +++---- ...naged-metadata-column-naming-scheme.golden | 16 ++++ .../output/tableflow/topic/update-help.golden | 19 +++-- ...naged-metadata-column-naming-scheme.golden | 21 +++++ test/tableflow_test.go | 2 + test/test-server/tableflow_handlers.go | 4 + 12 files changed, 155 insertions(+), 87 deletions(-) create mode 100644 test/fixtures/output/tableflow/topic/enable-topic-managed-metadata-column-naming-scheme.golden create mode 100644 test/fixtures/output/tableflow/topic/update-topic-managed-metadata-column-naming-scheme.golden diff --git a/go.mod b/go.mod index 949f24507d..dc63bf9b05 100644 --- a/go.mod +++ b/go.mod @@ -51,7 +51,7 @@ require ( github.com/confluentinc/ccloud-sdk-go-v2/service-quota v0.2.0 github.com/confluentinc/ccloud-sdk-go-v2/srcm v0.7.3 github.com/confluentinc/ccloud-sdk-go-v2/sso v0.0.1 - github.com/confluentinc/ccloud-sdk-go-v2/tableflow v0.6.0 + github.com/confluentinc/ccloud-sdk-go-v2/tableflow v0.7.0 github.com/confluentinc/ccloud-sdk-go-v2/usm v0.1.0 github.com/confluentinc/cmf-sdk-go v0.0.8 github.com/confluentinc/confluent-kafka-go/v2 v2.14.2 @@ -117,7 +117,7 @@ require ( go.uber.org/mock v0.4.0 golang.org/x/crypto v0.54.0 golang.org/x/exp v0.0.0-20260410095643-746e56fc9e2f - golang.org/x/oauth2 v0.35.0 + golang.org/x/oauth2 v0.36.0 golang.org/x/term v0.45.0 golang.org/x/text v0.40.0 google.golang.org/grpc v1.80.0 diff --git a/go.sum b/go.sum index 165d12317d..1814489dd5 100644 --- a/go.sum +++ b/go.sum @@ -244,8 +244,8 @@ github.com/confluentinc/ccloud-sdk-go-v2/srcm v0.7.3 h1:ozdDSJHruQIgtxS5hwz8Rp8p github.com/confluentinc/ccloud-sdk-go-v2/srcm v0.7.3/go.mod h1:cD0AeCMBAWBesmWxWCMgVYNABYgHJ/ahCj7b4HP2R2I= github.com/confluentinc/ccloud-sdk-go-v2/sso v0.0.1 h1:WZJYfgXJrvTIYQpCFps/qHF7T8ekgPlX/SFqx4EY2zQ= github.com/confluentinc/ccloud-sdk-go-v2/sso v0.0.1/go.mod h1:kB+MXWYYg9ohrTCb27LlfpTbuexAzyYAmum105ow0ho= -github.com/confluentinc/ccloud-sdk-go-v2/tableflow v0.6.0 h1:wrmpI4UJgWZ4rX1EYqUxUQrfsKMRDehQsIWjcMU/bzs= -github.com/confluentinc/ccloud-sdk-go-v2/tableflow v0.6.0/go.mod h1:myRmhUEWzpwGqWvdsNk79QH41pkre1G21vmySyGQiWA= +github.com/confluentinc/ccloud-sdk-go-v2/tableflow v0.7.0 h1:yf/1TiSCiAT0cD1dBZB6SU+hUAAL+uCezLx5MM9aHCA= +github.com/confluentinc/ccloud-sdk-go-v2/tableflow v0.7.0/go.mod h1:A0gPt6EjmzCV+CftXu2TeefIOiha2OTbATVOPp0ccx0= github.com/confluentinc/ccloud-sdk-go-v2/usm v0.1.0 h1:rF9cKecDCowq+oDWjf8rSpXXZHAnVXowIsT/OXF4MOI= github.com/confluentinc/ccloud-sdk-go-v2/usm v0.1.0/go.mod h1:umhEDvQp/5h0ALKBpYTQOmFwaWrvilnbE8Rkzh6oJ4Q= github.com/confluentinc/cmf-sdk-go v0.0.8 h1:ziiV4/lNKh49m7RxwPCDL04OFV54lA3SkeynsQ2SDC8= @@ -847,8 +847,8 @@ golang.org/x/oauth2 v0.0.0-20201109201403-9fd604954f58/go.mod h1:KelEdhl1UZF7XfJ golang.org/x/oauth2 v0.0.0-20201208152858-08078c50e5b5/go.mod h1:KelEdhl1UZF7XfJ4dDtk6s++YSgaE7mD/BuKKDLBl4A= golang.org/x/oauth2 v0.0.0-20210218202405-ba52d332ba99/go.mod h1:KelEdhl1UZF7XfJ4dDtk6s++YSgaE7mD/BuKKDLBl4A= golang.org/x/oauth2 v0.0.0-20210323180902-22b0adad7558/go.mod h1:KelEdhl1UZF7XfJ4dDtk6s++YSgaE7mD/BuKKDLBl4A= -golang.org/x/oauth2 v0.35.0 h1:Mv2mzuHuZuY2+bkyWXIHMfhNdJAdwW3FuWeCPYN5GVQ= -golang.org/x/oauth2 v0.35.0/go.mod h1:lzm5WQJQwKZ3nwavOZ3IS5Aulzxi68dUSgRHujetwEA= +golang.org/x/oauth2 v0.36.0 h1:peZ/1z27fi9hUOFCAZaHyrpWG5lwe0RJEEEeH0ThlIs= +golang.org/x/oauth2 v0.36.0/go.mod h1:YDBUJMTkDnJS+A4BP4eZBjCqtokkg1hODuPjwiGPO7Q= golang.org/x/sync v0.0.0-20180314180146-1d60e4601c6f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= diff --git a/internal/tableflow/command_topic.go b/internal/tableflow/command_topic.go index e460af9735..90242e897e 100644 --- a/internal/tableflow/command_topic.go +++ b/internal/tableflow/command_topic.go @@ -24,30 +24,31 @@ const ( ) type topicOut struct { - KafkaCluster string `human:"Kafka Cluster" serialized:"kafka_cluster"` - TopicName string `human:"Topic Name" serialized:"topic_name"` - EnableCompaction bool `human:"Enable Compaction" serialized:"enable_compaction"` - EnablePartitioning bool `human:"Enable Partitioning" serialized:"enable_partitioning"` - Environment string `human:"Environment" serialized:"environment"` - RecordFailureStrategy string `human:"Record Failure Strategy" serialized:"record_failure_strategy"` - ErrorHandling string `human:"Error Handling,omitempty" serialized:"error_handling,omitempty"` - LogTarget string `human:"Log Target,omitempty" serialized:"log_target,omitempty"` - RetentionMs string `human:"Retention Ms" serialized:"retention_ms"` - StorageType string `human:"Storage Type" serialized:"storage_type"` - ProviderIntegrationId string `human:"Provider Integration ID,omitempty" serialized:"provider_integration_id,omitempty"` - BucketName string `human:"Bucket Name,omitempty" serialized:"bucket_name,omitempty"` - BucketRegion string `human:"Bucket Region,omitempty" serialized:"bucket_region,omitempty"` - ContainerName string `human:"Container Name,omitempty" serialized:"container_name,omitempty"` - StorageAccountName string `human:"Storage Account Name,omitempty" serialized:"storage_account_name,omitempty"` - StorageRegion string `human:"Storage Region,omitempty" serialized:"storage_region ,omitempty"` - Suspended bool `human:"Suspended" serialized:"suspended"` - TableFormats string `human:"Table Formats" serialized:"table_formats"` - TablePath string `human:"Table Path" serialized:"table_path"` - Phase string `human:"Phase" serialized:"phase"` - CatalogSyncStatus map[string]string `human:"Catalog Sync Status,omitempty" serialized:"catalog_sync_status,omitempty"` - FailingTableFormat map[string]string `human:"Failing Table Format,omitempty" serialized:"failing_table_format,omitempty"` - ErrorMessage string `human:"Error Message,omitempty" serialized:"error_message,omitempty"` - WriteMode string `human:"Write Mode,omitempty" serialized:"write_mode,omitempty"` + KafkaCluster string `human:"Kafka Cluster" serialized:"kafka_cluster"` + TopicName string `human:"Topic Name" serialized:"topic_name"` + EnableCompaction bool `human:"Enable Compaction" serialized:"enable_compaction"` + EnablePartitioning bool `human:"Enable Partitioning" serialized:"enable_partitioning"` + Environment string `human:"Environment" serialized:"environment"` + RecordFailureStrategy string `human:"Record Failure Strategy" serialized:"record_failure_strategy"` + MetadataColumnNamingScheme string `human:"Metadata Column Naming Scheme,omitempty" serialized:"metadata_column_naming_scheme,omitempty"` + ErrorHandling string `human:"Error Handling,omitempty" serialized:"error_handling,omitempty"` + LogTarget string `human:"Log Target,omitempty" serialized:"log_target,omitempty"` + RetentionMs string `human:"Retention Ms" serialized:"retention_ms"` + StorageType string `human:"Storage Type" serialized:"storage_type"` + ProviderIntegrationId string `human:"Provider Integration ID,omitempty" serialized:"provider_integration_id,omitempty"` + BucketName string `human:"Bucket Name,omitempty" serialized:"bucket_name,omitempty"` + BucketRegion string `human:"Bucket Region,omitempty" serialized:"bucket_region,omitempty"` + ContainerName string `human:"Container Name,omitempty" serialized:"container_name,omitempty"` + StorageAccountName string `human:"Storage Account Name,omitempty" serialized:"storage_account_name,omitempty"` + StorageRegion string `human:"Storage Region,omitempty" serialized:"storage_region ,omitempty"` + Suspended bool `human:"Suspended" serialized:"suspended"` + TableFormats string `human:"Table Formats" serialized:"table_formats"` + TablePath string `human:"Table Path" serialized:"table_path"` + Phase string `human:"Phase" serialized:"phase"` + CatalogSyncStatus map[string]string `human:"Catalog Sync Status,omitempty" serialized:"catalog_sync_status,omitempty"` + FailingTableFormat map[string]string `human:"Failing Table Format,omitempty" serialized:"failing_table_format,omitempty"` + ErrorMessage string `human:"Error Message,omitempty" serialized:"error_message,omitempty"` + WriteMode string `human:"Write Mode,omitempty" serialized:"write_mode,omitempty"` } func (c *command) newTopicCommand() *cobra.Command { @@ -203,23 +204,24 @@ func printTopicTable(cmd *cobra.Command, topic tableflowv1.TableflowV1TableflowT strFormats := getFailingTableFormats(topic.Status.GetFailingTableFormats()) out := &topicOut{ - KafkaCluster: topic.GetSpec().KafkaCluster.GetId(), - TopicName: topic.Spec.GetDisplayName(), - EnableCompaction: topic.GetSpec().Config.GetEnableCompaction(), // should be read-only & true - EnablePartitioning: topic.GetSpec().Config.GetEnablePartitioning(), // should be read-only & true - TableFormats: strings.Join(topic.Spec.GetTableFormats(), ", "), - Environment: topic.GetSpec().Environment.GetId(), - RetentionMs: topic.GetSpec().Config.GetRetentionMs(), - RecordFailureStrategy: topic.GetSpec().Config.GetRecordFailureStrategy(), - ErrorHandling: getErrorHandlingMode(topic), - LogTarget: topic.GetSpec().Config.GetErrorHandling().TableflowV1ErrorHandlingLog.GetTarget(), // this Get function will return empty string if the ErrorHandling is not LOG - StorageType: storageType, - Suspended: topic.Spec.GetSuspended(), - Phase: topic.Status.GetPhase(), - CatalogSyncStatus: strStatus, - FailingTableFormat: strFormats, - ErrorMessage: topic.Status.GetErrorMessage(), - WriteMode: topic.Status.GetWriteMode(), + KafkaCluster: topic.GetSpec().KafkaCluster.GetId(), + TopicName: topic.Spec.GetDisplayName(), + EnableCompaction: topic.GetSpec().Config.GetEnableCompaction(), // should be read-only & true + EnablePartitioning: topic.GetSpec().Config.GetEnablePartitioning(), // should be read-only & true + TableFormats: strings.Join(topic.Spec.GetTableFormats(), ", "), + Environment: topic.GetSpec().Environment.GetId(), + RetentionMs: topic.GetSpec().Config.GetRetentionMs(), + RecordFailureStrategy: topic.GetSpec().Config.GetRecordFailureStrategy(), + MetadataColumnNamingScheme: topic.GetSpec().Config.GetMetadataColumnNamingScheme(), + ErrorHandling: getErrorHandlingMode(topic), + LogTarget: topic.GetSpec().Config.GetErrorHandling().TableflowV1ErrorHandlingLog.GetTarget(), // this Get function will return empty string if the ErrorHandling is not LOG + StorageType: storageType, + Suspended: topic.Spec.GetSuspended(), + Phase: topic.Status.GetPhase(), + CatalogSyncStatus: strStatus, + FailingTableFormat: strFormats, + ErrorMessage: topic.Status.GetErrorMessage(), + WriteMode: topic.Status.GetWriteMode(), } if storageType == byos { diff --git a/internal/tableflow/command_topic_enable.go b/internal/tableflow/command_topic_enable.go index 98af9bbcb0..542237068d 100644 --- a/internal/tableflow/command_topic_enable.go +++ b/internal/tableflow/command_topic_enable.go @@ -39,6 +39,7 @@ func (c *command) newTopicEnableCommand() *cobra.Command { cmd.Flags().String("provider-integration", "", "Specify the provider integration id.") cmd.Flags().String("bucket-name", "", "Specify the name of the AWS S3 bucket.") cmd.Flags().String("table-formats", "ICEBERG", "Specify the table formats, one of DELTA or ICEBERG.") + cmd.Flags().String("metadata-column-naming-scheme", "", "Specify the naming scheme for Tableflow's internal metadata columns in the materialized table, one of DEFAULT or PORTABLE.") cmd.Flags().String("storage-account-name", "", "Specify the storage account name for Azure Data Lake.") cmd.Flags().String("container-name", "", "Specify the container name for Azure Data Lake.") addErrorHandlingFlags(cmd) @@ -76,6 +77,11 @@ func (c *command) enable(cmd *cobra.Command, args []string) error { return err } + metadataColumnNamingScheme, err := cmd.Flags().GetString("metadata-column-naming-scheme") + if err != nil { + return err + } + errorHandling, err := cmd.Flags().GetString("error-handling") if err != nil { return err @@ -135,6 +141,10 @@ func (c *command) enable(cmd *cobra.Command, args []string) error { createTopic.Spec.Config.SetRecordFailureStrategy(recordFailureStrategy) } + if cmd.Flags().Changed("metadata-column-naming-scheme") { + createTopic.Spec.Config.SetMetadataColumnNamingScheme(metadataColumnNamingScheme) + } + if cmd.Flags().Changed("error-handling") { if strings.ToUpper(errorHandling) == suspend { createTopic.Spec.Config.ErrorHandling = &tableflowv1.TableflowV1TableFlowTopicConfigsSpecErrorHandlingOneOf{ diff --git a/internal/tableflow/command_topic_list.go b/internal/tableflow/command_topic_list.go index 4c21aba288..12ea31f17a 100644 --- a/internal/tableflow/command_topic_list.go +++ b/internal/tableflow/command_topic_list.go @@ -54,23 +54,24 @@ func (c *command) list(cmd *cobra.Command, _ []string) error { strFormats := getFailingTableFormats(topic.Status.GetFailingTableFormats()) out := &topicOut{ - KafkaCluster: topic.GetSpec().KafkaCluster.GetId(), - TopicName: topic.Spec.GetDisplayName(), - EnableCompaction: topic.GetSpec().Config.GetEnableCompaction(), // should be read-only & true - EnablePartitioning: topic.GetSpec().Config.GetEnablePartitioning(), // should be read-only & true - TableFormats: strings.Join(topic.Spec.GetTableFormats(), ", "), - Environment: topic.GetSpec().Environment.GetId(), - RetentionMs: topic.GetSpec().Config.GetRetentionMs(), - RecordFailureStrategy: topic.GetSpec().Config.GetRecordFailureStrategy(), - ErrorHandling: getErrorHandlingMode(topic), - LogTarget: topic.GetSpec().Config.GetErrorHandling().TableflowV1ErrorHandlingLog.GetTarget(), // this Get function will return empty string if the ErrorHandling is not LOG - StorageType: storageType, - Suspended: topic.Spec.GetSuspended(), - Phase: topic.Status.GetPhase(), - CatalogSyncStatus: strStatus, - FailingTableFormat: strFormats, - ErrorMessage: topic.Status.GetErrorMessage(), - WriteMode: topic.Status.GetWriteMode(), + KafkaCluster: topic.GetSpec().KafkaCluster.GetId(), + TopicName: topic.Spec.GetDisplayName(), + EnableCompaction: topic.GetSpec().Config.GetEnableCompaction(), // should be read-only & true + EnablePartitioning: topic.GetSpec().Config.GetEnablePartitioning(), // should be read-only & true + TableFormats: strings.Join(topic.Spec.GetTableFormats(), ", "), + Environment: topic.GetSpec().Environment.GetId(), + RetentionMs: topic.GetSpec().Config.GetRetentionMs(), + RecordFailureStrategy: topic.GetSpec().Config.GetRecordFailureStrategy(), + MetadataColumnNamingScheme: topic.GetSpec().Config.GetMetadataColumnNamingScheme(), + ErrorHandling: getErrorHandlingMode(topic), + LogTarget: topic.GetSpec().Config.GetErrorHandling().TableflowV1ErrorHandlingLog.GetTarget(), // this Get function will return empty string if the ErrorHandling is not LOG + StorageType: storageType, + Suspended: topic.Spec.GetSuspended(), + Phase: topic.Status.GetPhase(), + CatalogSyncStatus: strStatus, + FailingTableFormat: strFormats, + ErrorMessage: topic.Status.GetErrorMessage(), + WriteMode: topic.Status.GetWriteMode(), } if storageType == byos { diff --git a/internal/tableflow/command_topic_update.go b/internal/tableflow/command_topic_update.go index 6b396ab243..ecc125c87d 100644 --- a/internal/tableflow/command_topic_update.go +++ b/internal/tableflow/command_topic_update.go @@ -32,6 +32,7 @@ func (c *command) newTopicUpdateCommand() *cobra.Command { cmd.Flags().String("retention-ms", "", "Specify the Tableflow table retention time in milliseconds.") cmd.Flags().String("table-formats", "", "Specify the table formats, one of DELTA or ICEBERG.") + cmd.Flags().String("metadata-column-naming-scheme", "", "Specify the naming scheme for Tableflow's internal metadata columns in the materialized table, one of DEFAULT or PORTABLE.") addErrorHandlingFlags(cmd) pcmd.AddContextFlag(cmd, c.CLICommand) @@ -71,6 +72,11 @@ func (c *command) update(cmd *cobra.Command, args []string) error { return err } + metadataColumnNamingScheme, err := cmd.Flags().GetString("metadata-column-naming-scheme") + if err != nil { + return err + } + errorHandling, err := cmd.Flags().GetString("error-handling") if err != nil { return err @@ -101,6 +107,10 @@ func (c *command) update(cmd *cobra.Command, args []string) error { topicUpdate.Spec.Config.SetRecordFailureStrategy(recordFailureStrategy) } + if cmd.Flags().Changed("metadata-column-naming-scheme") { + topicUpdate.Spec.Config.SetMetadataColumnNamingScheme(metadataColumnNamingScheme) + } + if cmd.Flags().Changed("error-handling") { if strings.ToUpper(errorHandling) == suspend { topicUpdate.Spec.Config.ErrorHandling = &tableflowv1.TableflowV1TableFlowTopicConfigsSpecErrorHandlingOneOf{ diff --git a/test/fixtures/output/tableflow/topic/enable-help.golden b/test/fixtures/output/tableflow/topic/enable-help.golden index cc67edd4e7..bd67e143f4 100644 --- a/test/fixtures/output/tableflow/topic/enable-help.golden +++ b/test/fixtures/output/tableflow/topic/enable-help.golden @@ -16,20 +16,21 @@ Enable a confluent managed Tableflow topic related to a Kafka cluster. $ confluent tableflow topic enable my-tableflow-topic --cluster lkc-123456 --retention-ms 604800000 --storage-type MANAGED Flags: - --cluster string Kafka cluster ID. - --retention-ms string Specify the max age of snapshots (Iceberg) or versions (Delta) (snapshot/version expiration) to keep on the table in milliseconds for the Tableflow enabled topic. (default "604800000") - --storage-type string Specify the storage type of the Kafka cluster, one of MANAGED, BYOS or AzureDataLakeStorageGen2. (default "MANAGED") - --provider-integration string Specify the provider integration id. - --bucket-name string Specify the name of the AWS S3 bucket. - --table-formats string Specify the table formats, one of DELTA or ICEBERG. (default "ICEBERG") - --storage-account-name string Specify the storage account name for Azure Data Lake. - --container-name string Specify the container name for Azure Data Lake. - --error-handling string Specify the error handling strategy, one of SUSPEND, SKIP, or LOG. - --log-target string Specify the target topic for the LOG error handling strategy. - --context string CLI context name. - --environment string Environment ID. - -o, --output string Specify the output format as "human", "json", or "yaml". (default "human") - --record-failure-strategy string DEPRECATED: Specify the record failure strategy, one of SUSPEND or SKIP. + --cluster string Kafka cluster ID. + --retention-ms string Specify the max age of snapshots (Iceberg) or versions (Delta) (snapshot/version expiration) to keep on the table in milliseconds for the Tableflow enabled topic. (default "604800000") + --storage-type string Specify the storage type of the Kafka cluster, one of MANAGED, BYOS or AzureDataLakeStorageGen2. (default "MANAGED") + --provider-integration string Specify the provider integration id. + --bucket-name string Specify the name of the AWS S3 bucket. + --table-formats string Specify the table formats, one of DELTA or ICEBERG. (default "ICEBERG") + --metadata-column-naming-scheme string Specify the naming scheme for Tableflow's internal metadata columns in the materialized table, one of DEFAULT or PORTABLE. + --storage-account-name string Specify the storage account name for Azure Data Lake. + --container-name string Specify the container name for Azure Data Lake. + --error-handling string Specify the error handling strategy, one of SUSPEND, SKIP, or LOG. + --log-target string Specify the target topic for the LOG error handling strategy. + --context string CLI context name. + --environment string Environment ID. + -o, --output string Specify the output format as "human", "json", or "yaml". (default "human") + --record-failure-strategy string DEPRECATED: Specify the record failure strategy, one of SUSPEND or SKIP. Global Flags: -h, --help Show help for this command. diff --git a/test/fixtures/output/tableflow/topic/enable-topic-managed-metadata-column-naming-scheme.golden b/test/fixtures/output/tableflow/topic/enable-topic-managed-metadata-column-naming-scheme.golden new file mode 100644 index 0000000000..e5367ffae4 --- /dev/null +++ b/test/fixtures/output/tableflow/topic/enable-topic-managed-metadata-column-naming-scheme.golden @@ -0,0 +1,16 @@ ++-------------------------------+--------------------------------------------------------------------------+ +| Kafka Cluster | lkc-123456 | +| Topic Name | topic-managed | +| Enable Compaction | true | +| Enable Partitioning | false | +| Environment | | +| Record Failure Strategy | SUSPEND | +| Metadata Column Naming Scheme | PORTABLE | +| Retention Ms | 604800000 | +| Storage Type | MANAGED | +| Suspended | false | +| Table Formats | ICEBERG | +| Table Path | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | +| Phase | RUNNING | +| Write Mode | APPEND | ++-------------------------------+--------------------------------------------------------------------------+ diff --git a/test/fixtures/output/tableflow/topic/update-help.golden b/test/fixtures/output/tableflow/topic/update-help.golden index ca83726a37..6e45f73840 100644 --- a/test/fixtures/output/tableflow/topic/update-help.golden +++ b/test/fixtures/output/tableflow/topic/update-help.golden @@ -9,15 +9,16 @@ Update the refresh interval or retention time of Tableflow topic "my-tableflow-t $ confluent tableflow topic update my-tableflow-topic --cluster lkc-123456 --retention-ms 432000000 Flags: - --cluster string Kafka cluster ID. - --retention-ms string Specify the Tableflow table retention time in milliseconds. - --table-formats string Specify the table formats, one of DELTA or ICEBERG. - --error-handling string Specify the error handling strategy, one of SUSPEND, SKIP, or LOG. - --log-target string Specify the target topic for the LOG error handling strategy. - --context string CLI context name. - --environment string Environment ID. - -o, --output string Specify the output format as "human", "json", or "yaml". (default "human") - --record-failure-strategy string DEPRECATED: Specify the record failure strategy, one of SUSPEND or SKIP. + --cluster string Kafka cluster ID. + --retention-ms string Specify the Tableflow table retention time in milliseconds. + --table-formats string Specify the table formats, one of DELTA or ICEBERG. + --metadata-column-naming-scheme string Specify the naming scheme for Tableflow's internal metadata columns in the materialized table, one of DEFAULT or PORTABLE. + --error-handling string Specify the error handling strategy, one of SUSPEND, SKIP, or LOG. + --log-target string Specify the target topic for the LOG error handling strategy. + --context string CLI context name. + --environment string Environment ID. + -o, --output string Specify the output format as "human", "json", or "yaml". (default "human") + --record-failure-strategy string DEPRECATED: Specify the record failure strategy, one of SUSPEND or SKIP. Global Flags: -h, --help Show help for this command. diff --git a/test/fixtures/output/tableflow/topic/update-topic-managed-metadata-column-naming-scheme.golden b/test/fixtures/output/tableflow/topic/update-topic-managed-metadata-column-naming-scheme.golden new file mode 100644 index 0000000000..cf8676080e --- /dev/null +++ b/test/fixtures/output/tableflow/topic/update-topic-managed-metadata-column-naming-scheme.golden @@ -0,0 +1,21 @@ ++-------------------------------+--------------------------------------------------------------------------+ +| Kafka Cluster | lkc-123456 | +| Topic Name | topic-managed | +| Enable Compaction | true | +| Enable Partitioning | true | +| Environment | env-596 | +| Record Failure Strategy | SUSPEND | +| Metadata Column Naming Scheme | PORTABLE | +| Error Handling | SUSPEND | +| Retention Ms | 604800000 | +| Storage Type | MANAGED | +| Suspended | false | +| Table Formats | DELTA | +| Table Path | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | +| Phase | RUNNING | +| Catalog Sync Status | cat-id-123=SUCCESS | +| | cat-id-456=FAILED-Connection timeout | +| Failing Table Format | DELTA=Connection timeout | +| | ICEBERG=Schema validation failed | +| Write Mode | APPEND | ++-------------------------------+--------------------------------------------------------------------------+ diff --git a/test/tableflow_test.go b/test/tableflow_test.go index 529303a9a1..8526740357 100644 --- a/test/tableflow_test.go +++ b/test/tableflow_test.go @@ -63,10 +63,12 @@ func (s *CLITestSuite) TestTableflowTopic() { {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --error-handling SUSPEND", fixture: "tableflow/topic/enable-topic-managed-error-handling-suspend.golden"}, {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --error-handling SKIP", fixture: "tableflow/topic/enable-topic-managed-error-handling-skip.golden"}, {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --error-handling LOG --log-target log_topic", fixture: "tableflow/topic/enable-topic-managed-error-handling-log.golden"}, + {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --metadata-column-naming-scheme PORTABLE", fixture: "tableflow/topic/enable-topic-managed-metadata-column-naming-scheme.golden"}, {args: "tableflow topic update topic-byob --cluster lkc-123456 --retention-ms 432000000", fixture: "tableflow/topic/update-topic.golden"}, {args: "tableflow topic update topic-managed --cluster lkc-123456 --error-handling SUSPEND", fixture: "tableflow/topic/update-topic-managed-error-handling-suspend.golden"}, {args: "tableflow topic update topic-managed --cluster lkc-123456 --error-handling SKIP", fixture: "tableflow/topic/update-topic-managed-error-handling-skip.golden"}, {args: "tableflow topic update topic-managed --cluster lkc-123456 --error-handling LOG --log-target log_topic", fixture: "tableflow/topic/update-topic-managed-error-handling-log.golden"}, + {args: "tableflow topic update topic-managed --cluster lkc-123456 --metadata-column-naming-scheme PORTABLE", fixture: "tableflow/topic/update-topic-managed-metadata-column-naming-scheme.golden"}, {args: "tableflow topic update topic-managed --cluster lkc-123456 --log-target log_topic", fixture: "tableflow/topic/update-topic-managed-no-change.golden"}, {args: "tableflow topic update topic-error-log --cluster lkc-123456 --log-target log_topic", fixture: "tableflow/topic/update-topic-error-log.golden"}, {args: "tableflow topic disable topic-managed --cluster lkc-123456", input: "y\n", fixture: "tableflow/topic/disable-topic.golden"}, diff --git a/test/test-server/tableflow_handlers.go b/test/test-server/tableflow_handlers.go index 4d4b1a2459..b02d0e9a4a 100644 --- a/test/test-server/tableflow_handlers.go +++ b/test/test-server/tableflow_handlers.go @@ -165,6 +165,10 @@ func handleTableflowTopicUpdate(t *testing.T, display_name string) http.HandlerF tableflowTopic.Spec.Config.SetRetentionMs(body.Spec.Config.GetRetentionMs()) } + if body.Spec.Config.GetMetadataColumnNamingScheme() != "" { + tableflowTopic.Spec.Config.SetMetadataColumnNamingScheme(body.Spec.Config.GetMetadataColumnNamingScheme()) + } + if body.Spec.Config.HasErrorHandling() { tableflowTopic.Spec.Config.SetErrorHandling(body.Spec.Config.GetErrorHandling()) } From fb1d886880255aeb7baf032d287a05975f031cbc Mon Sep 17 00:00:00 2001 From: rajdangwal Date: Wed, 22 Jul 2026 18:55:10 +0530 Subject: [PATCH 2/4] USTORAGE-27286: Normalize scheme flag case; fix storage_region JSON key Address review feedback: - Uppercase-normalize --metadata-column-naming-scheme before sending, so values like "portable" work (consistent with --storage-type/--error-handling). The field is x-extensible-enum, so we deliberately don't hard-reject unknown values client-side; the server validates and returns a clear error. - Fix the pre-existing trailing space in the storage_region serialized tag ("storage_region ,omitempty"), which produced a malformed JSON/YAML key; regen the affected JSON goldens. - Regenerate list-topic.golden to include the Metadata Column Naming Scheme column (list output shows all columns). --- internal/tableflow/command_topic.go | 2 +- internal/tableflow/command_topic_enable.go | 2 +- internal/tableflow/command_topic_update.go | 2 +- .../topic/describe-topic-azure-json.golden | 2 +- .../output/tableflow/topic/list-topic-json.golden | 2 +- .../output/tableflow/topic/list-topic.golden | 14 +++++++------- test/tableflow_test.go | 2 +- 7 files changed, 13 insertions(+), 13 deletions(-) diff --git a/internal/tableflow/command_topic.go b/internal/tableflow/command_topic.go index 90242e897e..9e291ca874 100644 --- a/internal/tableflow/command_topic.go +++ b/internal/tableflow/command_topic.go @@ -40,7 +40,7 @@ type topicOut struct { BucketRegion string `human:"Bucket Region,omitempty" serialized:"bucket_region,omitempty"` ContainerName string `human:"Container Name,omitempty" serialized:"container_name,omitempty"` StorageAccountName string `human:"Storage Account Name,omitempty" serialized:"storage_account_name,omitempty"` - StorageRegion string `human:"Storage Region,omitempty" serialized:"storage_region ,omitempty"` + StorageRegion string `human:"Storage Region,omitempty" serialized:"storage_region,omitempty"` Suspended bool `human:"Suspended" serialized:"suspended"` TableFormats string `human:"Table Formats" serialized:"table_formats"` TablePath string `human:"Table Path" serialized:"table_path"` diff --git a/internal/tableflow/command_topic_enable.go b/internal/tableflow/command_topic_enable.go index 542237068d..3ffa8adcfa 100644 --- a/internal/tableflow/command_topic_enable.go +++ b/internal/tableflow/command_topic_enable.go @@ -142,7 +142,7 @@ func (c *command) enable(cmd *cobra.Command, args []string) error { } if cmd.Flags().Changed("metadata-column-naming-scheme") { - createTopic.Spec.Config.SetMetadataColumnNamingScheme(metadataColumnNamingScheme) + createTopic.Spec.Config.SetMetadataColumnNamingScheme(strings.ToUpper(metadataColumnNamingScheme)) } if cmd.Flags().Changed("error-handling") { diff --git a/internal/tableflow/command_topic_update.go b/internal/tableflow/command_topic_update.go index ecc125c87d..a9cd976f13 100644 --- a/internal/tableflow/command_topic_update.go +++ b/internal/tableflow/command_topic_update.go @@ -108,7 +108,7 @@ func (c *command) update(cmd *cobra.Command, args []string) error { } if cmd.Flags().Changed("metadata-column-naming-scheme") { - topicUpdate.Spec.Config.SetMetadataColumnNamingScheme(metadataColumnNamingScheme) + topicUpdate.Spec.Config.SetMetadataColumnNamingScheme(strings.ToUpper(metadataColumnNamingScheme)) } if cmd.Flags().Changed("error-handling") { diff --git a/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden b/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden index bc4638a910..2e2137a7d1 100644 --- a/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden +++ b/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden @@ -10,7 +10,7 @@ "provider_integration_id": "cspi-stgce89r7", "container_name": "Container1", "storage_account_name": "Acc1", - "storage_region ": "US1", + "storage_region": "US1", "suspended": false, "table_formats": "ICEBERG", "table_path": "s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2", diff --git a/test/fixtures/output/tableflow/topic/list-topic-json.golden b/test/fixtures/output/tableflow/topic/list-topic-json.golden index 35bea5e412..fe63359037 100644 --- a/test/fixtures/output/tableflow/topic/list-topic-json.golden +++ b/test/fixtures/output/tableflow/topic/list-topic-json.golden @@ -11,7 +11,7 @@ "provider_integration_id": "cspi-stgce89r7", "container_name": "Container1", "storage_account_name": "Acc1", - "storage_region ": "US1", + "storage_region": "US1", "suspended": false, "table_formats": "ICEBERG", "table_path": "s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2", diff --git a/test/fixtures/output/tableflow/topic/list-topic.golden b/test/fixtures/output/tableflow/topic/list-topic.golden index 53000e9f60..4d5d792ec7 100644 --- a/test/fixtures/output/tableflow/topic/list-topic.golden +++ b/test/fixtures/output/tableflow/topic/list-topic.golden @@ -1,7 +1,7 @@ - Kafka Cluster | Topic Name | Enable Compaction | Enable Partitioning | Environment | Record Failure Strategy | Error Handling | Log Target | Retention Ms | Storage Type | Provider Integration ID | Bucket Name | Bucket Region | Container Name | Storage Account Name | Storage Region | Suspended | Table Formats | Table Path | Phase | Catalog Sync Status | Failing Table Format | Error Message | Write Mode -----------------+---------------+-------------------+---------------------+-------------+-------------------------+----------------+------------+--------------+--------------------------+-------------------------+-------------+---------------+----------------+----------------------+----------------+-----------+---------------+---------------------------------------------------------------------------+---------+---------------------+----------------------------------+---------------+------------- - lkc-123456 | topic-azure | true | true | env-596 | SKIP | | | 604800000 | AzureDataLakeStorageGen2 | cspi-stgce89r7 | | | Container1 | Acc1 | US1 | false | ICEBERG | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2 | RUNNING | | | | UPSERT - lkc-123456 | topic-byob | true | true | env-596 | SKIP | SKIP | | 604800000 | BYOS | cspi-stgce89r7 | bucket_1 | us-east-1 | | | | false | ICEBERG | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | RUNNING | cat-id-123=SUCCESS | DELTA=Connection timeout | | UPSERT - | | | | | | | | | | | | | | | | | | | | cat-id-456=FAILED | ICEBERG=Schema validation failed | | - lkc-123456 | topic-managed | true | true | env-596 | SUSPEND | SUSPEND | | 604800000 | MANAGED | | | | | | | false | DELTA | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | RUNNING | cat-id-123=SUCCESS | DELTA=Connection timeout | | APPEND - | | | | | | | | | | | | | | | | | | | | cat-id-456=FAILED | ICEBERG=Schema validation failed | | + Kafka Cluster | Topic Name | Enable Compaction | Enable Partitioning | Environment | Record Failure Strategy | Metadata Column Naming Scheme | Error Handling | Log Target | Retention Ms | Storage Type | Provider Integration ID | Bucket Name | Bucket Region | Container Name | Storage Account Name | Storage Region | Suspended | Table Formats | Table Path | Phase | Catalog Sync Status | Failing Table Format | Error Message | Write Mode +----------------+---------------+-------------------+---------------------+-------------+-------------------------+-------------------------------+----------------+------------+--------------+--------------------------+-------------------------+-------------+---------------+----------------+----------------------+----------------+-----------+---------------+---------------------------------------------------------------------------+---------+---------------------+----------------------------------+---------------+------------- + lkc-123456 | topic-azure | true | true | env-596 | SKIP | | | | 604800000 | AzureDataLakeStorageGen2 | cspi-stgce89r7 | | | Container1 | Acc1 | US1 | false | ICEBERG | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2 | RUNNING | | | | UPSERT + lkc-123456 | topic-byob | true | true | env-596 | SKIP | | SKIP | | 604800000 | BYOS | cspi-stgce89r7 | bucket_1 | us-east-1 | | | | false | ICEBERG | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | RUNNING | cat-id-123=SUCCESS | DELTA=Connection timeout | | UPSERT + | | | | | | | | | | | | | | | | | | | | | cat-id-456=FAILED | ICEBERG=Schema validation failed | | + lkc-123456 | topic-managed | true | true | env-596 | SUSPEND | | SUSPEND | | 604800000 | MANAGED | | | | | | | false | DELTA | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | RUNNING | cat-id-123=SUCCESS | DELTA=Connection timeout | | APPEND + | | | | | | | | | | | | | | | | | | | | | cat-id-456=FAILED | ICEBERG=Schema validation failed | | diff --git a/test/tableflow_test.go b/test/tableflow_test.go index 8526740357..bce8e26548 100644 --- a/test/tableflow_test.go +++ b/test/tableflow_test.go @@ -63,7 +63,7 @@ func (s *CLITestSuite) TestTableflowTopic() { {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --error-handling SUSPEND", fixture: "tableflow/topic/enable-topic-managed-error-handling-suspend.golden"}, {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --error-handling SKIP", fixture: "tableflow/topic/enable-topic-managed-error-handling-skip.golden"}, {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --error-handling LOG --log-target log_topic", fixture: "tableflow/topic/enable-topic-managed-error-handling-log.golden"}, - {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --metadata-column-naming-scheme PORTABLE", fixture: "tableflow/topic/enable-topic-managed-metadata-column-naming-scheme.golden"}, + {args: "tableflow topic enable topic-managed --cluster lkc-123456 --storage-type MANAGED --metadata-column-naming-scheme portable", fixture: "tableflow/topic/enable-topic-managed-metadata-column-naming-scheme.golden"}, {args: "tableflow topic update topic-byob --cluster lkc-123456 --retention-ms 432000000", fixture: "tableflow/topic/update-topic.golden"}, {args: "tableflow topic update topic-managed --cluster lkc-123456 --error-handling SUSPEND", fixture: "tableflow/topic/update-topic-managed-error-handling-suspend.golden"}, {args: "tableflow topic update topic-managed --cluster lkc-123456 --error-handling SKIP", fixture: "tableflow/topic/update-topic-managed-error-handling-skip.golden"}, From 76f03a9aaacc0633f80cdb7b3e79faabd400b212 Mon Sep 17 00:00:00 2001 From: rajdangwal Date: Wed, 22 Jul 2026 19:09:12 +0530 Subject: [PATCH 3/4] USTORAGE-27286: Allowlist metadata-column-naming-scheme in the flag-name linter The flag name exceeds the default 20-char / 2-delimiter limits. Add it to both ExcludeFlag allowlists (length and delimiter), matching the sibling record-failure-strategy flag, so the name stays aligned with the API/SDK field metadata_column_naming_scheme rather than abbreviating it. --- cmd/lint/main.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/cmd/lint/main.go b/cmd/lint/main.go index 4622615f84..d5aa628187 100644 --- a/cmd/lint/main.go +++ b/cmd/lint/main.go @@ -113,6 +113,7 @@ var flagRules = []linter.FlagRule{ "include-parent-scopes", "max-partition-memory-bytes", "message-send-max-retries", + "metadata-column-naming-scheme", "private-link-access-point", "record-failure-strategy", "remote-api-key", @@ -142,6 +143,7 @@ var flagRules = []linter.FlagRule{ "confluent-platform-kafka-cluster", "max-partition-memory-bytes", "message-send-max-retries", + "metadata-column-naming-scheme", "private-link-access-point", "require-crl-on-client-certificate", "schema-registry-api-key", From 40814db65106635e1b4121f892e70a9d34fbbfb1 Mon Sep 17 00:00:00 2001 From: rajdangwal Date: Wed, 22 Jul 2026 19:53:38 +0530 Subject: [PATCH 4/4] USTORAGE-27286: Revert unrelated storage_region fix; add read-back test coverage Revert the storage_region serialized-tag change: renaming a serialized field is a breaking change per AGENTS.md and is unrelated to this feature, so it shouldn't ride along in a minor PR. Set metadata_column_naming_scheme on the Azure mock topic so the describe and list goldens exercise a non-empty value read back from the API (previously only the enable/update echo paths covered it, leaving command_topic_list.go's read untested). --- internal/tableflow/command_topic.go | 2 +- .../topic/describe-topic-azure-json.golden | 3 +- .../topic/describe-topic-azure.golden | 39 ++++++++++--------- .../tableflow/topic/list-topic-json.golden | 3 +- .../output/tableflow/topic/list-topic.golden | 2 +- test/test-server/tableflow_handlers.go | 9 +++-- 6 files changed, 31 insertions(+), 27 deletions(-) diff --git a/internal/tableflow/command_topic.go b/internal/tableflow/command_topic.go index 9e291ca874..90242e897e 100644 --- a/internal/tableflow/command_topic.go +++ b/internal/tableflow/command_topic.go @@ -40,7 +40,7 @@ type topicOut struct { BucketRegion string `human:"Bucket Region,omitempty" serialized:"bucket_region,omitempty"` ContainerName string `human:"Container Name,omitempty" serialized:"container_name,omitempty"` StorageAccountName string `human:"Storage Account Name,omitempty" serialized:"storage_account_name,omitempty"` - StorageRegion string `human:"Storage Region,omitempty" serialized:"storage_region,omitempty"` + StorageRegion string `human:"Storage Region,omitempty" serialized:"storage_region ,omitempty"` Suspended bool `human:"Suspended" serialized:"suspended"` TableFormats string `human:"Table Formats" serialized:"table_formats"` TablePath string `human:"Table Path" serialized:"table_path"` diff --git a/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden b/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden index 2e2137a7d1..94ebfabbc2 100644 --- a/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden +++ b/test/fixtures/output/tableflow/topic/describe-topic-azure-json.golden @@ -5,12 +5,13 @@ "enable_partitioning": true, "environment": "env-596", "record_failure_strategy": "SKIP", + "metadata_column_naming_scheme": "PORTABLE", "retention_ms": "604800000", "storage_type": "AzureDataLakeStorageGen2", "provider_integration_id": "cspi-stgce89r7", "container_name": "Container1", "storage_account_name": "Acc1", - "storage_region": "US1", + "storage_region ": "US1", "suspended": false, "table_formats": "ICEBERG", "table_path": "s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2", diff --git a/test/fixtures/output/tableflow/topic/describe-topic-azure.golden b/test/fixtures/output/tableflow/topic/describe-topic-azure.golden index 27e74908d5..fb4b5844a3 100644 --- a/test/fixtures/output/tableflow/topic/describe-topic-azure.golden +++ b/test/fixtures/output/tableflow/topic/describe-topic-azure.golden @@ -1,19 +1,20 @@ -+-------------------------+---------------------------------------------------------------------------+ -| Kafka Cluster | lkc-123456 | -| Topic Name | topic-azure | -| Enable Compaction | true | -| Enable Partitioning | true | -| Environment | env-596 | -| Record Failure Strategy | SKIP | -| Retention Ms | 604800000 | -| Storage Type | AzureDataLakeStorageGen2 | -| Provider Integration ID | cspi-stgce89r7 | -| Container Name | Container1 | -| Storage Account Name | Acc1 | -| Storage Region | US1 | -| Suspended | false | -| Table Formats | ICEBERG | -| Table Path | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2 | -| Phase | RUNNING | -| Write Mode | UPSERT | -+-------------------------+---------------------------------------------------------------------------+ ++-------------------------------+---------------------------------------------------------------------------+ +| Kafka Cluster | lkc-123456 | +| Topic Name | topic-azure | +| Enable Compaction | true | +| Enable Partitioning | true | +| Environment | env-596 | +| Record Failure Strategy | SKIP | +| Metadata Column Naming Scheme | PORTABLE | +| Retention Ms | 604800000 | +| Storage Type | AzureDataLakeStorageGen2 | +| Provider Integration ID | cspi-stgce89r7 | +| Container Name | Container1 | +| Storage Account Name | Acc1 | +| Storage Region | US1 | +| Suspended | false | +| Table Formats | ICEBERG | +| Table Path | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2 | +| Phase | RUNNING | +| Write Mode | UPSERT | ++-------------------------------+---------------------------------------------------------------------------+ diff --git a/test/fixtures/output/tableflow/topic/list-topic-json.golden b/test/fixtures/output/tableflow/topic/list-topic-json.golden index fe63359037..082a16b200 100644 --- a/test/fixtures/output/tableflow/topic/list-topic-json.golden +++ b/test/fixtures/output/tableflow/topic/list-topic-json.golden @@ -6,12 +6,13 @@ "enable_partitioning": true, "environment": "env-596", "record_failure_strategy": "SKIP", + "metadata_column_naming_scheme": "PORTABLE", "retention_ms": "604800000", "storage_type": "AzureDataLakeStorageGen2", "provider_integration_id": "cspi-stgce89r7", "container_name": "Container1", "storage_account_name": "Acc1", - "storage_region": "US1", + "storage_region ": "US1", "suspended": false, "table_formats": "ICEBERG", "table_path": "s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2", diff --git a/test/fixtures/output/tableflow/topic/list-topic.golden b/test/fixtures/output/tableflow/topic/list-topic.golden index 4d5d792ec7..c3c6ce4dc1 100644 --- a/test/fixtures/output/tableflow/topic/list-topic.golden +++ b/test/fixtures/output/tableflow/topic/list-topic.golden @@ -1,6 +1,6 @@ Kafka Cluster | Topic Name | Enable Compaction | Enable Partitioning | Environment | Record Failure Strategy | Metadata Column Naming Scheme | Error Handling | Log Target | Retention Ms | Storage Type | Provider Integration ID | Bucket Name | Bucket Region | Container Name | Storage Account Name | Storage Region | Suspended | Table Formats | Table Path | Phase | Catalog Sync Status | Failing Table Format | Error Message | Write Mode ----------------+---------------+-------------------+---------------------+-------------+-------------------------+-------------------------------+----------------+------------+--------------+--------------------------+-------------------------+-------------+---------------+----------------+----------------------+----------------+-----------+---------------+---------------------------------------------------------------------------+---------+---------------------+----------------------------------+---------------+------------- - lkc-123456 | topic-azure | true | true | env-596 | SKIP | | | | 604800000 | AzureDataLakeStorageGen2 | cspi-stgce89r7 | | | Container1 | Acc1 | US1 | false | ICEBERG | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2 | RUNNING | | | | UPSERT + lkc-123456 | topic-azure | true | true | env-596 | SKIP | PORTABLE | | | 604800000 | AzureDataLakeStorageGen2 | cspi-stgce89r7 | | | Container1 | Acc1 | US1 | false | ICEBERG | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId2 | RUNNING | | | | UPSERT lkc-123456 | topic-byob | true | true | env-596 | SKIP | | SKIP | | 604800000 | BYOS | cspi-stgce89r7 | bucket_1 | us-east-1 | | | | false | ICEBERG | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | RUNNING | cat-id-123=SUCCESS | DELTA=Connection timeout | | UPSERT | | | | | | | | | | | | | | | | | | | | | cat-id-456=FAILED | ICEBERG=Schema validation failed | | lkc-123456 | topic-managed | true | true | env-596 | SUSPEND | | SUSPEND | | 604800000 | MANAGED | | | | | | | false | DELTA | s3://dummy-bucket-name-1//10011010/11101100/org-1/env-2/lkc-3/v1/tableId | RUNNING | cat-id-123=SUCCESS | DELTA=Connection timeout | | APPEND diff --git a/test/test-server/tableflow_handlers.go b/test/test-server/tableflow_handlers.go index b02d0e9a4a..3c526f8053 100644 --- a/test/test-server/tableflow_handlers.go +++ b/test/test-server/tableflow_handlers.go @@ -322,10 +322,11 @@ func getTopicAzure(display_name, environmentId, clusterId string) tableflowv1.Ta }, }, Config: &tableflowv1.TableflowV1TableFlowTopicConfigsSpec{ - EnableCompaction: tableflowv1.PtrBool(true), - EnablePartitioning: tableflowv1.PtrBool(true), // read-only property that needs confirmation, assuming constantly true for now - RetentionMs: tableflowv1.PtrString("604800000"), // 7 days to milliseconds - RecordFailureStrategy: tableflowv1.PtrString("SKIP"), + EnableCompaction: tableflowv1.PtrBool(true), + EnablePartitioning: tableflowv1.PtrBool(true), // read-only property that needs confirmation, assuming constantly true for now + RetentionMs: tableflowv1.PtrString("604800000"), // 7 days to milliseconds + RecordFailureStrategy: tableflowv1.PtrString("SKIP"), + MetadataColumnNamingScheme: tableflowv1.PtrString("PORTABLE"), }, TableFormats: &[]string{"ICEBERG"}, Environment: &tableflowv1.GlobalObjectReference{Id: environmentId},