Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
76 changes: 65 additions & 11 deletions internal/flink/command_environment_create.go
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
package flink

import (
"bytes"
"encoding/json"
"fmt"
"io"
"os"
"path/filepath"
"strings"
Expand All @@ -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 <name>",
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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 {
Expand All @@ -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
}
Comment on lines +172 to +199

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The comment makes sense, we can absorb it and make it part of the validation function as I suggest above.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Handled — the strict decoders now reject trailing content.


func jsonMarshalHelper(v interface{}, label string) (string, error) {
data, err := json.Marshal(v)
if err != nil {
Expand Down
9 changes: 1 addition & 8 deletions internal/flink/command_environment_update.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
config-overrides:
key: value
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
detached:
flinkConfiguration:
key1: value1
interactive:
flinkConfiguration:
key2: value2
Original file line number Diff line number Diff line change
@@ -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":{...}}}.
Original file line number Diff line number Diff line change
@@ -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":{...}}}.
Original file line number Diff line number Diff line change
@@ -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":{...}}}.
Original file line number Diff line number Diff line change
@@ -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"
}
}
}
}
Original file line number Diff line number Diff line change
@@ -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":{...}}}.
Original file line number Diff line number Diff line change
@@ -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":{...}}}.
6 changes: 6 additions & 0 deletions test/flink_onprem_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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" +
Expand All @@ -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" +
Expand Down