From 2391d663608ddc3a71c9c4145cc057d93c198e21 Mon Sep 17 00:00:00 2001 From: Paras Negi Date: Mon, 6 Jul 2026 09:34:39 +0530 Subject: [PATCH 1/2] CF-3931 : Reject mis-shaped --statement-defaults instead of silently ignoring it --- internal/flink/command_environment_create.go | 26 ++++++++++++++++--- internal/flink/command_environment_update.go | 2 +- .../environment/create-help-onprem.golden | 2 +- .../environment/create-no-namespace.golden | 2 +- .../create-statement-defaults-invalid.golden | 1 + .../environment/missing-flag-failure.golden | 2 +- .../environment/update-help-onprem.golden | 2 +- .../update-statement-defaults-invalid.golden | 1 + test/flink_onprem_test.go | 2 ++ 9 files changed, 31 insertions(+), 9 deletions(-) create mode 100644 test/fixtures/output/flink/environment/create-statement-defaults-invalid.golden create mode 100644 test/fixtures/output/flink/environment/update-statement-defaults-invalid.golden diff --git a/internal/flink/command_environment_create.go b/internal/flink/command_environment_create.go index e7e0f82ec9..44bfce6dc0 100644 --- a/internal/flink/command_environment_create.go +++ b/internal/flink/command_environment_create.go @@ -1,6 +1,7 @@ package flink import ( + "bytes" "encoding/json" "fmt" "os" @@ -26,7 +27,7 @@ func (c *command) newEnvironmentCreateCommand() *cobra.Command { cmd.Flags().String("kubernetes-namespace", "", "Kubernetes namespace to deploy Flink applications to.") cmd.Flags().String("defaults", "", "JSON string defining the environment's Flink application defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension).") - cmd.Flags().String("statement-defaults", "", "JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension).") + cmd.Flags().String("statement-defaults", "", `JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). Expected shape: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}.`) cmd.Flags().String("compute-pool-defaults", "", "JSON string defining the environment's Flink compute pool defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension).") addCmfFlagSet(cmd) @@ -124,18 +125,18 @@ func parseDefaultsAsGenericType[T any](input, label string) (T, error) { if err != nil { return out, fmt.Errorf("failed to read %s defaults JSON file: %w", label, err) } - err = json.Unmarshal(data, &out) + err = decodeStrictJson(data, &out) case ".yaml", ".yml": data, err = os.ReadFile(input) if err != nil { return out, fmt.Errorf("failed to read %s defaults YAML file: %w", label, err) } - err = yaml.Unmarshal(data, &out) + err = decodeStrictYaml(data, &out) default: // inline JSON string - err = json.Unmarshal([]byte(input), &out) + err = decodeStrictJson([]byte(input), &out) } if err != nil { @@ -144,6 +145,23 @@ func parseDefaultsAsGenericType[T any](input, label string) (T, error) { return out, nil } +// decodeStrictJson decodes JSON into out, rejecting keys that do not map to a +// known field. This surfaces mis-shaped input (for example the wrong nesting for +// --statement-defaults) instead of silently dropping it. Decoding into a map is +// unaffected, since a map has no unknown fields. +func decodeStrictJson(data []byte, out any) error { + decoder := json.NewDecoder(bytes.NewReader(data)) + decoder.DisallowUnknownFields() + return decoder.Decode(out) +} + +// decodeStrictYaml is the YAML counterpart of decodeStrictJson. +func decodeStrictYaml(data []byte, out any) error { + decoder := yaml.NewDecoder(bytes.NewReader(data)) + decoder.KnownFields(true) + return decoder.Decode(out) +} + func jsonMarshalHelper(v interface{}, label string) (string, error) { data, err := json.Marshal(v) if err != nil { diff --git a/internal/flink/command_environment_update.go b/internal/flink/command_environment_update.go index 50721ec519..7c2ddfc29c 100644 --- a/internal/flink/command_environment_update.go +++ b/internal/flink/command_environment_update.go @@ -21,7 +21,7 @@ func (c *command) newEnvironmentUpdateCommand() *cobra.Command { addCmfFlagSet(cmd) cmd.Flags().String("defaults", "", "JSON string defining the environment's Flink application defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension).") - cmd.Flags().String("statement-defaults", "", "JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension).") + cmd.Flags().String("statement-defaults", "", `JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). Expected shape: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}.`) cmd.Flags().String("compute-pool-defaults", "", "JSON string defining the environment's Flink compute pool defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension).") pcmd.AddOutputFlag(cmd) diff --git a/test/fixtures/output/flink/environment/create-help-onprem.golden b/test/fixtures/output/flink/environment/create-help-onprem.golden index d55d420ed2..cf34270fdd 100644 --- a/test/fixtures/output/flink/environment/create-help-onprem.golden +++ b/test/fixtures/output/flink/environment/create-help-onprem.golden @@ -6,7 +6,7 @@ Usage: Flags: --kubernetes-namespace string REQUIRED: Kubernetes namespace to deploy Flink applications to. --defaults string JSON string defining the environment's Flink application defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). - --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). + --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). Expected shape: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. --compute-pool-defaults string JSON string defining the environment's Flink compute pool defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). --url string Base URL of the Confluent Manager for Apache Flink (CMF). Environment variable "CONFLUENT_CMF_URL" may be set in place of this flag. --client-key-path string Path to client private key for mTLS authentication. Environment variable "CONFLUENT_CMF_CLIENT_KEY_PATH" may be set in place of this flag. diff --git a/test/fixtures/output/flink/environment/create-no-namespace.golden b/test/fixtures/output/flink/environment/create-no-namespace.golden index 2e557bda41..088f641ffd 100644 --- a/test/fixtures/output/flink/environment/create-no-namespace.golden +++ b/test/fixtures/output/flink/environment/create-no-namespace.golden @@ -5,7 +5,7 @@ Usage: Flags: --kubernetes-namespace string REQUIRED: Kubernetes namespace to deploy Flink applications to. --defaults string JSON string defining the environment's Flink application defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). - --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). + --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). Expected shape: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. --compute-pool-defaults string JSON string defining the environment's Flink compute pool defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). --url string Base URL of the Confluent Manager for Apache Flink (CMF). Environment variable "CONFLUENT_CMF_URL" may be set in place of this flag. --client-key-path string Path to client private key for mTLS authentication. Environment variable "CONFLUENT_CMF_CLIENT_KEY_PATH" may be set in place of this flag. diff --git a/test/fixtures/output/flink/environment/create-statement-defaults-invalid.golden b/test/fixtures/output/flink/environment/create-statement-defaults-invalid.golden new file mode 100644 index 0000000000..80ce1f3003 --- /dev/null +++ b/test/fixtures/output/flink/environment/create-statement-defaults-invalid.golden @@ -0,0 +1 @@ +Error: failed to parse statement defaults: json: unknown field "config-overrides" diff --git a/test/fixtures/output/flink/environment/missing-flag-failure.golden b/test/fixtures/output/flink/environment/missing-flag-failure.golden index 6fede7f99e..b6bd070a27 100644 --- a/test/fixtures/output/flink/environment/missing-flag-failure.golden +++ b/test/fixtures/output/flink/environment/missing-flag-failure.golden @@ -8,7 +8,7 @@ Flags: --client-cert-path string Path to client cert to be verified by Confluent Manager for Apache Flink. Include for mTLS authentication. Environment variable "CONFLUENT_CMF_CLIENT_CERT_PATH" may be set in place of this flag. --certificate-authority-path string Path to a PEM-encoded Certificate Authority to verify the Confluent Manager for Apache Flink connection. Environment variable "CONFLUENT_CMF_CERTIFICATE_AUTHORITY_PATH" may be set in place of this flag. --defaults string JSON string defining the environment's Flink application defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). - --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). + --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). Expected shape: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. --compute-pool-defaults string JSON string defining the environment's Flink compute pool defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). -o, --output string Specify the output format as "human", "json", or "yaml". (default "human") diff --git a/test/fixtures/output/flink/environment/update-help-onprem.golden b/test/fixtures/output/flink/environment/update-help-onprem.golden index 1799958182..fb47a373f3 100644 --- a/test/fixtures/output/flink/environment/update-help-onprem.golden +++ b/test/fixtures/output/flink/environment/update-help-onprem.golden @@ -9,7 +9,7 @@ Flags: --client-cert-path string Path to client cert to be verified by Confluent Manager for Apache Flink. Include for mTLS authentication. Environment variable "CONFLUENT_CMF_CLIENT_CERT_PATH" may be set in place of this flag. --certificate-authority-path string Path to a PEM-encoded Certificate Authority to verify the Confluent Manager for Apache Flink connection. Environment variable "CONFLUENT_CMF_CERTIFICATE_AUTHORITY_PATH" may be set in place of this flag. --defaults string JSON string defining the environment's Flink application defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). - --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). + --statement-defaults string JSON string defining the environment's Flink statement defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). Expected shape: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. --compute-pool-defaults string JSON string defining the environment's Flink compute pool defaults, or path to a file to read defaults from (with .yml, .yaml or .json extension). -o, --output string Specify the output format as "human", "json", or "yaml". (default "human") diff --git a/test/fixtures/output/flink/environment/update-statement-defaults-invalid.golden b/test/fixtures/output/flink/environment/update-statement-defaults-invalid.golden new file mode 100644 index 0000000000..80ce1f3003 --- /dev/null +++ b/test/fixtures/output/flink/environment/update-statement-defaults-invalid.golden @@ -0,0 +1 @@ +Error: failed to parse statement defaults: json: unknown field "config-overrides" diff --git a/test/flink_onprem_test.go b/test/flink_onprem_test.go index c98ca8c4b6..e6b47d01b1 100644 --- a/test/flink_onprem_test.go +++ b/test/flink_onprem_test.go @@ -264,6 +264,7 @@ func (s *CLITestSuite) TestFlinkEnvironmentCreate() { {args: "flink environment create default-failure --kubernetes-namespace default-staging", fixture: "flink/environment/create-failure.golden", exitCode: 1}, {args: "flink environment create default --kubernetes-namespace default-staging", fixture: "flink/environment/create-existing.golden", exitCode: 1}, {args: "flink environment create default", fixture: "flink/environment/create-no-namespace.golden", exitCode: 1}, + {args: "flink environment create default-2 --kubernetes-namespace default-staging --statement-defaults '{\"config-overrides\":{\"key\":\"value\"}}'", fixture: "flink/environment/create-statement-defaults-invalid.golden", exitCode: 1}, // success with application, statement and compute pool defaults {args: "flink environment create default-2" + " --defaults test/fixtures/input/flink/environment/application-defaults.json" + @@ -288,6 +289,7 @@ func (s *CLITestSuite) TestFlinkEnvironmentUpdate() { {args: "flink environment update non-existent --defaults '{\"property\": \"value\"}'", fixture: "flink/environment/update-non-existent.golden", exitCode: 1}, {args: "flink environment update get-failure --defaults '{\"property\": \"value\"}'", fixture: "flink/environment/update-get-failure.golden", exitCode: 1}, {args: "flink environment update missing-flag-failure", fixture: "flink/environment/missing-flag-failure.golden", exitCode: 1}, + {args: "flink environment update default --statement-defaults '{\"config-overrides\":{\"key\":\"value\"}}'", fixture: "flink/environment/update-statement-defaults-invalid.golden", exitCode: 1}, // success with application, statement and compute pool defaults {args: "flink environment update default" + " --defaults test/fixtures/input/flink/environment/application-defaults.json" + From 599851acd50ee7a9420b7df4a887eed2aa7cdd85 Mon Sep 17 00:00:00 2001 From: Paras Negi Date: Tue, 7 Jul 2026 11:17:01 +0530 Subject: [PATCH 2/2] CF-3931 : Reject trailing data and extra documents in strict defaults decoding --- internal/flink/command_environment_create.go | 31 ++++++++++++++----- .../create-statement-defaults-trailing.golden | 1 + test/flink_onprem_test.go | 1 + 3 files changed, 26 insertions(+), 7 deletions(-) create mode 100644 test/fixtures/output/flink/environment/create-statement-defaults-trailing.golden diff --git a/internal/flink/command_environment_create.go b/internal/flink/command_environment_create.go index 44bfce6dc0..aa9019e17a 100644 --- a/internal/flink/command_environment_create.go +++ b/internal/flink/command_environment_create.go @@ -4,6 +4,7 @@ import ( "bytes" "encoding/json" "fmt" + "io" "os" "path/filepath" "strings" @@ -145,21 +146,37 @@ func parseDefaultsAsGenericType[T any](input, label string) (T, error) { return out, nil } -// decodeStrictJson decodes JSON into out, rejecting keys that do not map to a -// known field. This surfaces mis-shaped input (for example the wrong nesting for -// --statement-defaults) instead of silently dropping it. Decoding into a map is -// unaffected, since a map has no unknown fields. +// decodeStrictJson decodes a single JSON value into out, rejecting unknown +// fields and trailing data so mis-shaped input surfaces instead of being +// silently dropped. Decoding into a map is unaffected (a map has no unknown +// fields). func decodeStrictJson(data []byte, out any) error { decoder := json.NewDecoder(bytes.NewReader(data)) decoder.DisallowUnknownFields() - return decoder.Decode(out) + if err := decoder.Decode(out); err != nil { + return err + } + if decoder.More() { + return fmt.Errorf("unexpected trailing data after JSON value") + } + return nil } -// decodeStrictYaml is the YAML counterpart of decodeStrictJson. +// decodeStrictYaml is the YAML counterpart of decodeStrictJson; it also rejects +// unknown fields and any additional documents. func decodeStrictYaml(data []byte, out any) error { decoder := yaml.NewDecoder(bytes.NewReader(data)) decoder.KnownFields(true) - return decoder.Decode(out) + if err := decoder.Decode(out); err != nil { + return err + } + if err := decoder.Decode(&struct{}{}); err != io.EOF { + if err != nil { + return err + } + return fmt.Errorf("unexpected additional YAML document") + } + return nil } func jsonMarshalHelper(v interface{}, label string) (string, error) { diff --git a/test/fixtures/output/flink/environment/create-statement-defaults-trailing.golden b/test/fixtures/output/flink/environment/create-statement-defaults-trailing.golden new file mode 100644 index 0000000000..81adfcda1d --- /dev/null +++ b/test/fixtures/output/flink/environment/create-statement-defaults-trailing.golden @@ -0,0 +1 @@ +Error: failed to parse statement defaults: unexpected trailing data after JSON value diff --git a/test/flink_onprem_test.go b/test/flink_onprem_test.go index e6b47d01b1..29e329fc25 100644 --- a/test/flink_onprem_test.go +++ b/test/flink_onprem_test.go @@ -265,6 +265,7 @@ func (s *CLITestSuite) TestFlinkEnvironmentCreate() { {args: "flink environment create default --kubernetes-namespace default-staging", fixture: "flink/environment/create-existing.golden", exitCode: 1}, {args: "flink environment create default", fixture: "flink/environment/create-no-namespace.golden", exitCode: 1}, {args: "flink environment create default-2 --kubernetes-namespace default-staging --statement-defaults '{\"config-overrides\":{\"key\":\"value\"}}'", fixture: "flink/environment/create-statement-defaults-invalid.golden", exitCode: 1}, + {args: "flink environment create default-2 --kubernetes-namespace default-staging --statement-defaults '{\"interactive\":{}}{\"detached\":{}}'", fixture: "flink/environment/create-statement-defaults-trailing.golden", exitCode: 1}, // success with application, statement and compute pool defaults {args: "flink environment create default-2" + " --defaults test/fixtures/input/flink/environment/application-defaults.json" +