diff --git a/internal/flink/command_environment_create.go b/internal/flink/command_environment_create.go index e7e0f82ec9..d07c587360 100644 --- a/internal/flink/command_environment_create.go +++ b/internal/flink/command_environment_create.go @@ -1,8 +1,10 @@ package flink import ( + "bytes" "encoding/json" "fmt" + "io" "os" "path/filepath" "strings" @@ -13,9 +15,15 @@ import ( cmfsdk "github.com/confluentinc/cmf-sdk-go/v1" pcmd "github.com/confluentinc/cli/v4/pkg/cmd" + "github.com/confluentinc/cli/v4/pkg/errors" "github.com/confluentinc/cli/v4/pkg/output" ) +// statementDefaultsShapeSuggestion documents the expected `--statement-defaults` +// shape. It is surfaced as a suggestion when parsing fails rather than in the +// flag help so it stays next to the error the user actually hit. +const statementDefaultsShapeSuggestion = `Provide statement defaults matching the expected shape, for example: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}.` + func (c *command) newEnvironmentCreateCommand() *cobra.Command { cmd := &cobra.Command{ Use: "create ", @@ -80,16 +88,9 @@ func (c *command) environmentCreate(cmd *cobra.Command, args []string) error { } } if defaultsStatement != "" { - defaultsStatementParsedLocal, err := parseDefaultsAsGenericType[LocalAllStatementDefaults1](defaultsStatement, "statement") - if err != nil { + if defaultsStatementParsed, err = parseStatementDefaults(defaultsStatement); err != nil { return err } - if defaultsStatementParsedLocal.Detached != nil { - defaultsStatementParsed.SetDetached(cmfsdk.StatementDefaults{FlinkConfiguration: defaultsStatementParsedLocal.Detached.FlinkConfiguration}) - } - if defaultsStatementParsedLocal.Interactive != nil { - defaultsStatementParsed.SetInteractive(cmfsdk.StatementDefaults{FlinkConfiguration: defaultsStatementParsedLocal.Interactive.FlinkConfiguration}) - } } var postEnvironment cmfsdk.PostEnvironment @@ -112,6 +113,26 @@ func (c *command) environmentCreate(cmd *cobra.Command, args []string) error { return output.SerializedOutput(cmd, localEnv) } +// parseStatementDefaults strictly parses the `--statement-defaults` value into +// the SDK type, rejecting unknown/mis-nested fields. On failure it surfaces the +// expected shape as a suggestion instead of silently dropping the input. +func parseStatementDefaults(input string) (cmfsdk.AllStatementDefaults1, error) { + var statementDefaults cmfsdk.AllStatementDefaults1 + + parsed, err := parseDefaultsAsGenericType[LocalAllStatementDefaults1](input, "statement") + if err != nil { + return statementDefaults, errors.NewErrorWithSuggestions(err.Error(), statementDefaultsShapeSuggestion) + } + + if parsed.Detached != nil { + statementDefaults.SetDetached(cmfsdk.StatementDefaults{FlinkConfiguration: parsed.Detached.FlinkConfiguration}) + } + if parsed.Interactive != nil { + statementDefaults.SetInteractive(cmfsdk.StatementDefaults{FlinkConfiguration: parsed.Interactive.FlinkConfiguration}) + } + return statementDefaults, nil +} + func parseDefaultsAsGenericType[T any](input, label string) (T, error) { var out T var data []byte @@ -124,18 +145,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 +165,39 @@ func parseDefaultsAsGenericType[T any](input, label string) (T, error) { return out, nil } +// 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() + 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; 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) + 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) { 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..e43eea0a6d 100644 --- a/internal/flink/command_environment_update.go +++ b/internal/flink/command_environment_update.go @@ -69,16 +69,9 @@ func (c *command) environmentUpdate(cmd *cobra.Command, args []string) error { } } if defaultsStatement != "" { - defaultsStatementParsedLocal, err := parseDefaultsAsGenericType[LocalAllStatementDefaults1](defaultsStatement, "statement") - if err != nil { + if defaultsStatementParsed, err = parseStatementDefaults(defaultsStatement); err != nil { return err } - if defaultsStatementParsedLocal.Detached != nil { - defaultsStatementParsed.Detached = &cmfsdk.StatementDefaults{FlinkConfiguration: defaultsStatementParsedLocal.Detached.FlinkConfiguration} - } - if defaultsStatementParsedLocal.Interactive != nil { - defaultsStatementParsed.Interactive = &cmfsdk.StatementDefaults{FlinkConfiguration: defaultsStatementParsedLocal.Interactive.FlinkConfiguration} - } } var postEnvironment cmfsdk.PostEnvironment diff --git a/test/fixtures/input/flink/environment/statement-defaults-invalid.yaml b/test/fixtures/input/flink/environment/statement-defaults-invalid.yaml new file mode 100644 index 0000000000..cf6fd0d0f3 --- /dev/null +++ b/test/fixtures/input/flink/environment/statement-defaults-invalid.yaml @@ -0,0 +1,2 @@ +config-overrides: + key: value diff --git a/test/fixtures/input/flink/environment/statement-defaults.yaml b/test/fixtures/input/flink/environment/statement-defaults.yaml new file mode 100644 index 0000000000..55df3e6130 --- /dev/null +++ b/test/fixtures/input/flink/environment/statement-defaults.yaml @@ -0,0 +1,6 @@ +detached: + flinkConfiguration: + key1: value1 +interactive: + flinkConfiguration: + key2: value2 diff --git a/test/fixtures/output/flink/environment/create-statement-defaults-invalid-yaml.golden b/test/fixtures/output/flink/environment/create-statement-defaults-invalid-yaml.golden new file mode 100644 index 0000000000..14d9196a7a --- /dev/null +++ b/test/fixtures/output/flink/environment/create-statement-defaults-invalid-yaml.golden @@ -0,0 +1,5 @@ +Error: failed to parse statement defaults: yaml: unmarshal errors: + line 1: field config-overrides not found in type flink.LocalAllStatementDefaults1 + +Suggestions: + Provide statement defaults matching the expected shape, for example: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. 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..39677c45bd --- /dev/null +++ b/test/fixtures/output/flink/environment/create-statement-defaults-invalid.golden @@ -0,0 +1,4 @@ +Error: failed to parse statement defaults: json: unknown field "config-overrides" + +Suggestions: + Provide statement defaults matching the expected shape, for example: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. 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..8c8cfbd5d7 --- /dev/null +++ b/test/fixtures/output/flink/environment/create-statement-defaults-trailing.golden @@ -0,0 +1,4 @@ +Error: failed to parse statement defaults: unexpected trailing data after JSON value + +Suggestions: + Provide statement defaults matching the expected shape, for example: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. diff --git a/test/fixtures/output/flink/environment/create-statement-defaults-yaml-json.golden b/test/fixtures/output/flink/environment/create-statement-defaults-yaml-json.golden new file mode 100644 index 0000000000..24175138b3 --- /dev/null +++ b/test/fixtures/output/flink/environment/create-statement-defaults-yaml-json.golden @@ -0,0 +1,18 @@ +{ + "name": "default-2", + "created_time": "2024-09-10T23:00:00Z", + "updated_time": "2024-09-10T23:00:00Z", + "kubernetesNamespace": "default-staging", + "statementDefaults": { + "detached": { + "flinkConfiguration": { + "key1": "value1" + } + }, + "interactive": { + "flinkConfiguration": { + "key2": "value2" + } + } + } +} diff --git a/test/fixtures/output/flink/environment/update-statement-defaults-invalid-yaml.golden b/test/fixtures/output/flink/environment/update-statement-defaults-invalid-yaml.golden new file mode 100644 index 0000000000..14d9196a7a --- /dev/null +++ b/test/fixtures/output/flink/environment/update-statement-defaults-invalid-yaml.golden @@ -0,0 +1,5 @@ +Error: failed to parse statement defaults: yaml: unmarshal errors: + line 1: field config-overrides not found in type flink.LocalAllStatementDefaults1 + +Suggestions: + Provide statement defaults matching the expected shape, for example: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. 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..39677c45bd --- /dev/null +++ b/test/fixtures/output/flink/environment/update-statement-defaults-invalid.golden @@ -0,0 +1,4 @@ +Error: failed to parse statement defaults: json: unknown field "config-overrides" + +Suggestions: + Provide statement defaults matching the expected shape, for example: {"detached":{"flinkConfiguration":{...}},"interactive":{"flinkConfiguration":{...}}}. diff --git a/test/flink_onprem_test.go b/test/flink_onprem_test.go index c98ca8c4b6..426236491d 100644 --- a/test/flink_onprem_test.go +++ b/test/flink_onprem_test.go @@ -264,6 +264,10 @@ 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}, + {args: "flink environment create default-2 --kubernetes-namespace default-staging --statement-defaults '{\"interactive\":{}}{\"detached\":{}}'", fixture: "flink/environment/create-statement-defaults-trailing.golden", exitCode: 1}, + {args: "flink environment create default-2 --kubernetes-namespace default-staging --statement-defaults test/fixtures/input/flink/environment/statement-defaults-invalid.yaml", fixture: "flink/environment/create-statement-defaults-invalid-yaml.golden", exitCode: 1}, + {args: "flink environment create default-2 --kubernetes-namespace default-staging --statement-defaults test/fixtures/input/flink/environment/statement-defaults.yaml --output json", fixture: "flink/environment/create-statement-defaults-yaml-json.golden"}, // 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 +292,8 @@ 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}, + {args: "flink environment update default --statement-defaults test/fixtures/input/flink/environment/statement-defaults-invalid.yaml", fixture: "flink/environment/update-statement-defaults-invalid-yaml.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" +