diff --git a/internal/flink/command_application_list.go b/internal/flink/command_application_list.go index 63000b8001..4477cb217c 100644 --- a/internal/flink/command_application_list.go +++ b/internal/flink/command_application_list.go @@ -2,13 +2,22 @@ package flink import ( "fmt" + "slices" + "strings" "github.com/spf13/cobra" pcmd "github.com/confluentinc/cli/v4/pkg/cmd" "github.com/confluentinc/cli/v4/pkg/output" + "github.com/confluentinc/cli/v4/pkg/utils" ) +// allowedApplicationStatuses lists the Flink job states recognized by the CMF applications +// "state=" filter, per the cmf-sdk-go GetApplications filter documentation. Unknown values are +// still forwarded (the server returns no matches rather than erroring); this list only drives +// the advisory --status warning. +var allowedApplicationStatuses = []string{"RUNNING", "FINISHED", "FAILED", "CANCELED", "RECONCILING", "COMPLETED", "UNKNOWN"} + func (c *command) newApplicationListCommand() *cobra.Command { cmd := &cobra.Command{ Use: "list", @@ -18,6 +27,8 @@ func (c *command) newApplicationListCommand() *cobra.Command { } cmd.Flags().String("environment", "", "Name of the Flink environment.") + cmd.Flags().String("name", "", `Filter the Flink applications by name. Supports wildcards, for example "my-app*".`) + cmd.Flags().String("status", "", "Filter the Flink applications by status.") addPageSizeFlag(cmd) addCmfFlagSet(cmd) pcmd.AddOutputFlag(cmd) @@ -33,6 +44,22 @@ func (c *command) applicationList(cmd *cobra.Command, _ []string) error { return err } + name, err := cmd.Flags().GetString("name") + if err != nil { + return err + } + + status, err := cmd.Flags().GetString("status") + if err != nil { + return err + } + if status != "" { + status = strings.ToUpper(status) + if !slices.Contains(allowedApplicationStatuses, status) { + output.ErrPrintf(c.Config.EnableColor, "[WARN] Invalid status %q. Valid statuses are %s.\n", status, utils.ArrayToCommaDelimitedString(allowedApplicationStatuses, "and")) + } + } + pageSize, err := getPageSize(cmd) if err != nil { return err @@ -43,7 +70,7 @@ func (c *command) applicationList(cmd *cobra.Command, _ []string) error { return err } - applications, err := client.ListApplications(c.createContext(), environment, pageSize) + applications, err := client.ListApplications(c.createContext(), environment, buildApplicationFilter(name, status), pageSize) if err != nil { return err } @@ -82,3 +109,19 @@ func (c *command) applicationList(cmd *cobra.Command, _ []string) error { return output.SerializedOutput(cmd, localApps) } + +// buildApplicationFilter composes the CMF applications "filter" query from the user-facing +// --name and --status flags. The grammar (comma-separated "key=value" expressions, "name=" +// with an optional "*" suffix wildcard, and "state=" for status) follows the CMF applications +// list API. Values are not escaped: Kubernetes application names and Flink states cannot +// contain "," or "=", so no ambiguity arises. The caller normalizes and warns about status. +func buildApplicationFilter(name, status string) string { + filters := make([]string, 0, 2) + if name != "" { + filters = append(filters, fmt.Sprintf("name=%s", name)) + } + if status != "" { + filters = append(filters, fmt.Sprintf("state=%s", status)) + } + return strings.Join(filters, ",") +} diff --git a/internal/flink/command_application_list_test.go b/internal/flink/command_application_list_test.go new file mode 100644 index 0000000000..129811d396 --- /dev/null +++ b/internal/flink/command_application_list_test.go @@ -0,0 +1,28 @@ +package flink + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestBuildApplicationFilter(t *testing.T) { + tests := []struct { + name string + appn string + status string + want string + }{ + {name: "empty", appn: "", status: "", want: ""}, + {name: "name only", appn: "my-app", status: "", want: "name=my-app"}, + {name: "name wildcard", appn: "my-app*", status: "", want: "name=my-app*"}, + {name: "status only", appn: "", status: "RUNNING", want: "state=RUNNING"}, + {name: "name and status", appn: "a*", status: "RUNNING", want: "name=a*,state=RUNNING"}, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + require.Equal(t, test.want, buildApplicationFilter(test.appn, test.status)) + }) + } +} diff --git a/pkg/flink/cmf_rest_client.go b/pkg/flink/cmf_rest_client.go index 83a76a0672..b5856f1b27 100644 --- a/pkg/flink/cmf_rest_client.go +++ b/pkg/flink/cmf_rest_client.go @@ -189,9 +189,14 @@ func (cmfClient *CmfRestClient) DescribeApplication(ctx context.Context, environ return cmfApplication, nil } -func (cmfClient *CmfRestClient) ListApplications(ctx context.Context, environment string, pageSize int32) ([]cmfsdk.FlinkApplication, error) { +func (cmfClient *CmfRestClient) ListApplications(ctx context.Context, environment, filter string, pageSize int32) ([]cmfsdk.FlinkApplication, error) { + request := cmfClient.FlinkApplicationsApi.GetApplications(ctx, environment) + if filter != "" { + request = request.Filter(filter) + } + return listAllPages(pageSize, func(page, size int32) ([]cmfsdk.FlinkApplication, error) { - applicationsPage, httpResponse, err := cmfClient.FlinkApplicationsApi.GetApplications(ctx, environment).Page(page).Size(size).Execute() + applicationsPage, httpResponse, err := request.Page(page).Size(size).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list applications in the environment "%s": %s`, environment, parsedErr) } diff --git a/test/fixtures/output/flink/application/list-env-missing.golden b/test/fixtures/output/flink/application/list-env-missing.golden index 8358c03deb..878cdbfd0b 100644 --- a/test/fixtures/output/flink/application/list-env-missing.golden +++ b/test/fixtures/output/flink/application/list-env-missing.golden @@ -4,6 +4,8 @@ Usage: Flags: --environment string REQUIRED: Name of the Flink environment. + --name string Filter the Flink applications by name. Supports wildcards, for example "my-app*". + --status string Filter the Flink applications by status. --page-size int Number of results to fetch per API request. Defaults to 100. --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/application/list-help-onprem.golden b/test/fixtures/output/flink/application/list-help-onprem.golden index 750da75330..441d3ba4d3 100644 --- a/test/fixtures/output/flink/application/list-help-onprem.golden +++ b/test/fixtures/output/flink/application/list-help-onprem.golden @@ -5,6 +5,8 @@ Usage: Flags: --environment string REQUIRED: Name of the Flink environment. + --name string Filter the Flink applications by name. Supports wildcards, for example "my-app*". + --status string Filter the Flink applications by status. --page-size int Number of results to fetch per API request. Defaults to 100. --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/application/list-name-filter-json.golden b/test/fixtures/output/flink/application/list-name-filter-json.golden new file mode 100644 index 0000000000..273d16e858 --- /dev/null +++ b/test/fixtures/output/flink/application/list-name-filter-json.golden @@ -0,0 +1,83 @@ +[ + { + "apiVersion": "cmf.confluent.io/v1", + "kind": "FlinkApplication", + "metadata": { + "name": "default-application-s" + }, + "spec": { + "flinkConfiguration": { + "metrics.reporter.prom.factory.class": "org.apache.flink.metrics.prometheus.PrometheusReporterFactory", + "metrics.reporter.prom.port": "9249-9250", + "taskmanager.numberOfTaskSlots": "8" + }, + "flinkVersion": "v1_19", + "image": "confluentinc/cp-flink:1.19.1-cp1", + "job": { + "jarURI": "local:///opt/flink/examples/streaming/StateMachineExample.jar", + "parallelism": 3, + "state": "running", + "upgradeMode": "stateless" + }, + "jobManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + }, + "serviceAccount": "flink", + "taskManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + } + }, + "status": { + "clusterInfo": { + "flink-revision": "89d0b8f @ 2024-06-22T13:19:31+02:00", + "flink-version": "1.19.1-cp1", + "total-cpu": "3.0", + "total-memory": "3296722944" + }, + "error": null, + "jobManagerDeploymentStatus": "DEPLOYING", + "jobStatus": { + "checkpointInfo": { + "formatType": null, + "lastCheckpoint": null, + "lastPeriodicCheckpointTimestamp": 0, + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "jobId": "dcabb1ad6c40495bc2d7fa7a0097c5aa", + "jobName": "State machine job", + "savepointInfo": { + "formatType": null, + "lastPeriodicSavepointTimestamp": 0, + "lastSavepoint": null, + "savepointHistory": [], + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "startTime": "1726640263746", + "state": "RECONCILING", + "updateTime": "1726640280561" + }, + "lifecycleState": "DEPLOYED", + "observedGeneration": 4, + "reconciliationStatus": { + "lastReconciledSpec": "", + "lastStableSpec": "", + "reconciliationTimestamp": 1726640346899, + "state": "DEPLOYED" + }, + "taskManager": { + "labelSelector": "component=taskmanager,app=basic-example", + "replicas": 1 + } + } + } +] diff --git a/test/fixtures/output/flink/application/list-name-status-json.golden b/test/fixtures/output/flink/application/list-name-status-json.golden new file mode 100644 index 0000000000..230f62ffdc --- /dev/null +++ b/test/fixtures/output/flink/application/list-name-status-json.golden @@ -0,0 +1,83 @@ +[ + { + "apiVersion": "cmf.confluent.io/v1", + "kind": "FlinkApplication", + "metadata": { + "name": "default-application-1" + }, + "spec": { + "flinkConfiguration": { + "metrics.reporter.prom.factory.class": "org.apache.flink.metrics.prometheus.PrometheusReporterFactory", + "metrics.reporter.prom.port": "9249-9250", + "taskmanager.numberOfTaskSlots": "8" + }, + "flinkVersion": "v1_19", + "image": "confluentinc/cp-flink:1.19.1-cp1", + "job": { + "jarURI": "local:///opt/flink/examples/streaming/StateMachineExample.jar", + "parallelism": 3, + "state": "running", + "upgradeMode": "stateless" + }, + "jobManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + }, + "serviceAccount": "flink", + "taskManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + } + }, + "status": { + "clusterInfo": { + "flink-revision": "89d0b8f @ 2024-06-22T13:19:31+02:00", + "flink-version": "1.19.1-cp1", + "total-cpu": "3.0", + "total-memory": "3296722944" + }, + "error": null, + "jobManagerDeploymentStatus": "DEPLOYING", + "jobStatus": { + "checkpointInfo": { + "formatType": null, + "lastCheckpoint": null, + "lastPeriodicCheckpointTimestamp": 0, + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "jobId": "dcabb1ad6c40495bc2d7fa7a0097c5aa", + "jobName": "State machine job", + "savepointInfo": { + "formatType": null, + "lastPeriodicSavepointTimestamp": 0, + "lastSavepoint": null, + "savepointHistory": [], + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "startTime": "1726640263746", + "state": "RECONCILING", + "updateTime": "1726640280561" + }, + "lifecycleState": "DEPLOYED", + "observedGeneration": 4, + "reconciliationStatus": { + "lastReconciledSpec": "", + "lastStableSpec": "", + "reconciliationTimestamp": 1726640346899, + "state": "DEPLOYED" + }, + "taskManager": { + "labelSelector": "component=taskmanager,app=basic-example", + "replicas": 1 + } + } + } +] diff --git a/test/fixtures/output/flink/application/list-name-wildcard-json.golden b/test/fixtures/output/flink/application/list-name-wildcard-json.golden new file mode 100644 index 0000000000..230f62ffdc --- /dev/null +++ b/test/fixtures/output/flink/application/list-name-wildcard-json.golden @@ -0,0 +1,83 @@ +[ + { + "apiVersion": "cmf.confluent.io/v1", + "kind": "FlinkApplication", + "metadata": { + "name": "default-application-1" + }, + "spec": { + "flinkConfiguration": { + "metrics.reporter.prom.factory.class": "org.apache.flink.metrics.prometheus.PrometheusReporterFactory", + "metrics.reporter.prom.port": "9249-9250", + "taskmanager.numberOfTaskSlots": "8" + }, + "flinkVersion": "v1_19", + "image": "confluentinc/cp-flink:1.19.1-cp1", + "job": { + "jarURI": "local:///opt/flink/examples/streaming/StateMachineExample.jar", + "parallelism": 3, + "state": "running", + "upgradeMode": "stateless" + }, + "jobManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + }, + "serviceAccount": "flink", + "taskManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + } + }, + "status": { + "clusterInfo": { + "flink-revision": "89d0b8f @ 2024-06-22T13:19:31+02:00", + "flink-version": "1.19.1-cp1", + "total-cpu": "3.0", + "total-memory": "3296722944" + }, + "error": null, + "jobManagerDeploymentStatus": "DEPLOYING", + "jobStatus": { + "checkpointInfo": { + "formatType": null, + "lastCheckpoint": null, + "lastPeriodicCheckpointTimestamp": 0, + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "jobId": "dcabb1ad6c40495bc2d7fa7a0097c5aa", + "jobName": "State machine job", + "savepointInfo": { + "formatType": null, + "lastPeriodicSavepointTimestamp": 0, + "lastSavepoint": null, + "savepointHistory": [], + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "startTime": "1726640263746", + "state": "RECONCILING", + "updateTime": "1726640280561" + }, + "lifecycleState": "DEPLOYED", + "observedGeneration": 4, + "reconciliationStatus": { + "lastReconciledSpec": "", + "lastStableSpec": "", + "reconciliationTimestamp": 1726640346899, + "state": "DEPLOYED" + }, + "taskManager": { + "labelSelector": "component=taskmanager,app=basic-example", + "replicas": 1 + } + } + } +] diff --git a/test/fixtures/output/flink/application/list-status-filter-json.golden b/test/fixtures/output/flink/application/list-status-filter-json.golden new file mode 100644 index 0000000000..fee5108863 --- /dev/null +++ b/test/fixtures/output/flink/application/list-status-filter-json.golden @@ -0,0 +1,245 @@ +[ + { + "apiVersion": "cmf.confluent.io/v1", + "kind": "FlinkApplication", + "metadata": { + "name": "default-application-1" + }, + "spec": { + "flinkConfiguration": { + "metrics.reporter.prom.factory.class": "org.apache.flink.metrics.prometheus.PrometheusReporterFactory", + "metrics.reporter.prom.port": "9249-9250", + "taskmanager.numberOfTaskSlots": "8" + }, + "flinkVersion": "v1_19", + "image": "confluentinc/cp-flink:1.19.1-cp1", + "job": { + "jarURI": "local:///opt/flink/examples/streaming/StateMachineExample.jar", + "parallelism": 3, + "state": "running", + "upgradeMode": "stateless" + }, + "jobManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + }, + "serviceAccount": "flink", + "taskManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + } + }, + "status": { + "clusterInfo": { + "flink-revision": "89d0b8f @ 2024-06-22T13:19:31+02:00", + "flink-version": "1.19.1-cp1", + "total-cpu": "3.0", + "total-memory": "3296722944" + }, + "error": null, + "jobManagerDeploymentStatus": "DEPLOYING", + "jobStatus": { + "checkpointInfo": { + "formatType": null, + "lastCheckpoint": null, + "lastPeriodicCheckpointTimestamp": 0, + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "jobId": "dcabb1ad6c40495bc2d7fa7a0097c5aa", + "jobName": "State machine job", + "savepointInfo": { + "formatType": null, + "lastPeriodicSavepointTimestamp": 0, + "lastSavepoint": null, + "savepointHistory": [], + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "startTime": "1726640263746", + "state": "RECONCILING", + "updateTime": "1726640280561" + }, + "lifecycleState": "DEPLOYED", + "observedGeneration": 4, + "reconciliationStatus": { + "lastReconciledSpec": "", + "lastStableSpec": "", + "reconciliationTimestamp": 1726640346899, + "state": "DEPLOYED" + }, + "taskManager": { + "labelSelector": "component=taskmanager,app=basic-example", + "replicas": 1 + } + } + }, + { + "apiVersion": "cmf.confluent.io/v1", + "kind": "FlinkApplication", + "metadata": { + "name": "default-application-2" + }, + "spec": { + "flinkConfiguration": { + "metrics.reporter.prom.factory.class": "org.apache.flink.metrics.prometheus.PrometheusReporterFactory", + "metrics.reporter.prom.port": "9249-9250", + "taskmanager.numberOfTaskSlots": "8" + }, + "flinkVersion": "v1_19", + "image": "confluentinc/cp-flink:1.19.1-cp1", + "job": { + "jarURI": "local:///opt/flink/examples/streaming/StateMachineExample.jar", + "parallelism": 3, + "state": "running", + "upgradeMode": "stateless" + }, + "jobManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + }, + "serviceAccount": "flink", + "taskManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + } + }, + "status": { + "clusterInfo": { + "flink-revision": "89d0b8f @ 2024-06-22T13:19:31+02:00", + "flink-version": "1.19.1-cp1", + "total-cpu": "3.0", + "total-memory": "3296722944" + }, + "error": null, + "jobManagerDeploymentStatus": "DEPLOYING", + "jobStatus": { + "checkpointInfo": { + "formatType": null, + "lastCheckpoint": null, + "lastPeriodicCheckpointTimestamp": 0, + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "jobId": "dcabb1ad6c40495bc2d7fa7a0097c5aa", + "jobName": "State machine job", + "savepointInfo": { + "formatType": null, + "lastPeriodicSavepointTimestamp": 0, + "lastSavepoint": null, + "savepointHistory": [], + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "startTime": "1726640263746", + "state": "RECONCILING", + "updateTime": "1726640280561" + }, + "lifecycleState": "DEPLOYED", + "observedGeneration": 4, + "reconciliationStatus": { + "lastReconciledSpec": "", + "lastStableSpec": "", + "reconciliationTimestamp": 1726640346899, + "state": "DEPLOYED" + }, + "taskManager": { + "labelSelector": "component=taskmanager,app=basic-example", + "replicas": 1 + } + } + }, + { + "apiVersion": "cmf.confluent.io/v1", + "kind": "FlinkApplication", + "metadata": { + "name": "default-application-s" + }, + "spec": { + "flinkConfiguration": { + "metrics.reporter.prom.factory.class": "org.apache.flink.metrics.prometheus.PrometheusReporterFactory", + "metrics.reporter.prom.port": "9249-9250", + "taskmanager.numberOfTaskSlots": "8" + }, + "flinkVersion": "v1_19", + "image": "confluentinc/cp-flink:1.19.1-cp1", + "job": { + "jarURI": "local:///opt/flink/examples/streaming/StateMachineExample.jar", + "parallelism": 3, + "state": "running", + "upgradeMode": "stateless" + }, + "jobManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + }, + "serviceAccount": "flink", + "taskManager": { + "resource": { + "cpu": 1, + "memory": "1048m" + } + } + }, + "status": { + "clusterInfo": { + "flink-revision": "89d0b8f @ 2024-06-22T13:19:31+02:00", + "flink-version": "1.19.1-cp1", + "total-cpu": "3.0", + "total-memory": "3296722944" + }, + "error": null, + "jobManagerDeploymentStatus": "DEPLOYING", + "jobStatus": { + "checkpointInfo": { + "formatType": null, + "lastCheckpoint": null, + "lastPeriodicCheckpointTimestamp": 0, + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "jobId": "dcabb1ad6c40495bc2d7fa7a0097c5aa", + "jobName": "State machine job", + "savepointInfo": { + "formatType": null, + "lastPeriodicSavepointTimestamp": 0, + "lastSavepoint": null, + "savepointHistory": [], + "triggerId": null, + "triggerTimestamp": null, + "triggerType": null + }, + "startTime": "1726640263746", + "state": "RECONCILING", + "updateTime": "1726640280561" + }, + "lifecycleState": "DEPLOYED", + "observedGeneration": 4, + "reconciliationStatus": { + "lastReconciledSpec": "", + "lastStableSpec": "", + "reconciliationTimestamp": 1726640346899, + "state": "DEPLOYED" + }, + "taskManager": { + "labelSelector": "component=taskmanager,app=basic-example", + "replicas": 1 + } + } + } +] diff --git a/test/fixtures/output/flink/application/list-status-invalid-json.golden b/test/fixtures/output/flink/application/list-status-invalid-json.golden new file mode 100644 index 0000000000..73e46fa904 --- /dev/null +++ b/test/fixtures/output/flink/application/list-status-invalid-json.golden @@ -0,0 +1,2 @@ +[WARN] Invalid status "BOGUS". Valid statuses are "RUNNING", "FINISHED", "FAILED", "CANCELED", "RECONCILING", "COMPLETED", and "UNKNOWN". +[] diff --git a/test/fixtures/output/flink/application/list-status-no-match-json.golden b/test/fixtures/output/flink/application/list-status-no-match-json.golden new file mode 100644 index 0000000000..fe51488c70 --- /dev/null +++ b/test/fixtures/output/flink/application/list-status-no-match-json.golden @@ -0,0 +1 @@ +[] diff --git a/test/flink_onprem_test.go b/test/flink_onprem_test.go index c7643858c1..5f997f08d7 100644 --- a/test/flink_onprem_test.go +++ b/test/flink_onprem_test.go @@ -37,6 +37,13 @@ func (s *CLITestSuite) TestFlinkApplicationList() { {args: "flink application list --environment default --output human", fixture: "flink/application/list-human.golden"}, // pagination: a small page size still returns the full list, fetched across multiple round trips {args: "flink application list --environment default --page-size 2 --output json", fixture: "flink/application/list-json.golden"}, + // filtering + {args: "flink application list --environment default --name default-application-s --output json", fixture: "flink/application/list-name-filter-json.golden"}, + {args: "flink application list --environment default --name default-application-1* --output json", fixture: "flink/application/list-name-wildcard-json.golden"}, + {args: "flink application list --environment default --status reconciling --output json", fixture: "flink/application/list-status-filter-json.golden"}, + {args: "flink application list --environment default --status failed --output json", fixture: "flink/application/list-status-no-match-json.golden"}, + {args: "flink application list --environment default --status bogus --output json", fixture: "flink/application/list-status-invalid-json.golden"}, + {args: "flink application list --environment default --name default-application-1* --status reconciling --output json", fixture: "flink/application/list-name-status-json.golden"}, } runIntegrationTestsWithMultipleAuth(s, tests) diff --git a/test/test-server/flink_onprem_handler.go b/test/test-server/flink_onprem_handler.go index 655a600d78..4b488c14b3 100644 --- a/test/test-server/flink_onprem_handler.go +++ b/test/test-server/flink_onprem_handler.go @@ -146,6 +146,53 @@ func paginateApplications(all []cmfsdk.FlinkApplication, pageParam, sizeParam st return all[start:end] } +// filterApplications emulates the CMF server-side "filter" query param for the applications +// endpoint. It understands comma-separated "name=" (with optional "*" suffix wildcard) +// and "state=" expressions, as built by the CLI's --name/--status flags. +func filterApplications(items []cmfsdk.FlinkApplication, filter string) []cmfsdk.FlinkApplication { + if filter == "" { + return items + } + + for _, expr := range strings.Split(filter, ",") { + key, value, found := strings.Cut(expr, "=") + if !found { + continue + } + matched := make([]cmfsdk.FlinkApplication, 0, len(items)) + for _, item := range items { + if applicationMatchesFilter(item, key, value) { + matched = append(matched, item) + } + } + items = matched + } + return items +} + +func applicationMatchesFilter(app cmfsdk.FlinkApplication, key, value string) bool { + switch key { + case "name": + name, _ := app.Metadata["name"].(string) + if prefix, isWildcard := strings.CutSuffix(value, "*"); isWildcard { + return strings.HasPrefix(name, prefix) + } + return name == value + case "state": + if app.Status == nil { + return false + } + jobStatus, ok := (*app.Status)["jobStatus"].(map[string]interface{}) + if !ok { + return false + } + state, _ := jobStatus["state"].(string) + return strings.EqualFold(state, value) + default: + return true + } +} + // Helper function to create a Flink environment. func createEnvironment(name string, namespace string) cmfsdk.Environment { createdTime := time.Date(2024, time.September, 10, 23, 0, 0, 0, time.UTC) @@ -725,6 +772,7 @@ func handleCmfApplications(t *testing.T) http.HandlerFunc { allItems = []cmfsdk.FlinkApplication{createApplication("update-failure-application")} } + allItems = filterApplications(allItems, r.URL.Query().Get("filter")) applicationsPage := map[string]interface{}{ "items": paginateApplications(allItems, r.URL.Query().Get("page"), r.URL.Query().Get("size")), }