From a91ce7a4596e2583714b8267bf0f5be2a23365ba Mon Sep 17 00:00:00 2001 From: Paras Negi Date: Fri, 24 Jul 2026 16:43:39 +0530 Subject: [PATCH] [CF-4076] Add stable sort to paginated CMF list calls The on-prem (CMF) Flink list commands page through results with offset pagination but requested no sort order. With no ORDER BY, the server can return rows in a different order across separate page requests, so rows near a page boundary can be silently skipped or duplicated once a result set exceeds the page size of 100. Request a unique, stable sort on every paginated list call so offset paging is deterministic: - most resources sort by "name" (maps server-side to each unique id column) - secret mappings sort by "name","uid" ("name" has no unique constraint) - application events sort by "creationTimestamp,desc","name" (newest-first plus a unique tiebreaker) Add regression tests asserting a >100-item list returns every item exactly once and that each list call sends the expected sort key. Co-Authored-By: Claude Opus 4.8 --- pkg/flink/cmf_rest_client.go | 34 +++++--- pkg/flink/cmf_rest_client_test.go | 136 ++++++++++++++++++++++++++++++ 2 files changed, 157 insertions(+), 13 deletions(-) create mode 100644 pkg/flink/cmf_rest_client_test.go diff --git a/pkg/flink/cmf_rest_client.go b/pkg/flink/cmf_rest_client.go index 8518ad1423..497b0dddd6 100644 --- a/pkg/flink/cmf_rest_client.go +++ b/pkg/flink/cmf_rest_client.go @@ -45,6 +45,12 @@ type CmfRestClient struct { AuthToken string } +// sortByName gives paginated list calls a stable order so offset paging cannot +// silently skip or duplicate rows. Each list below is scoped (by environment, +// catalog, or parent), and within that scope "name" resolves to a unique column. +// Don't reuse it where "name" isn't unique in scope — see ListSecretMappings. +var sortByName = []string{"name"} + func NewCmfRestHttpClient(restFlags *OnPremCMFRestFlagValues) (*http.Client, error) { var err error httpClient := utils.DefaultClient() @@ -197,7 +203,7 @@ func (cmfClient *CmfRestClient) ListApplications(ctx context.Context, environmen done := false for !done { - applicationsPage, httpResponse, err := cmfClient.FlinkApplicationsApi.GetApplications(ctx, environment).Page(currentPageNumber).Size(pageSize).Execute() + applicationsPage, httpResponse, err := cmfClient.FlinkApplicationsApi.GetApplications(ctx, environment).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list applications in the environment "%s": %s`, environment, parsedErr) } @@ -237,7 +243,8 @@ func (cmfClient *CmfRestClient) ListApplicationEvents(ctx context.Context, envir done := false for !done { - eventsPage, httpResponse, err := cmfClient.FlinkApplicationsApi.GetApplicationEvents(ctx, environment, application).Page(currentPageNumber).Size(pageSize).Execute() + // Sort newest-first, with the unique "name" as a tiebreaker so paging is stable. + eventsPage, httpResponse, err := cmfClient.FlinkApplicationsApi.GetApplicationEvents(ctx, environment, application).Sort([]string{"creationTimestamp,desc", "name"}).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list events for application "%s" in the environment "%s": %s`, application, environment, parsedErr) } @@ -289,7 +296,7 @@ func (cmfClient *CmfRestClient) ListEnvironments(ctx context.Context) ([]cmfsdk. var currentPageNumber int32 = 0 for !done { - environmentsPage, httpResponse, err := cmfClient.EnvironmentsApi.GetEnvironments(ctx).Page(currentPageNumber).Size(pageSize).Execute() + environmentsPage, httpResponse, err := cmfClient.EnvironmentsApi.GetEnvironments(ctx).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf("failed to list environments: %s", parsedErr) } @@ -382,9 +389,9 @@ func (cmfClient *CmfRestClient) ListSavepoint(ctx context.Context, environment, var httpResponse *_nethttp.Response var err error if isStatement { - savepointsPage, httpResponse, err = cmfClient.SavepointsApi.GetSavepointsForFlinkStatement(ctx, environment, statement).Page(currentPageNumber).Size(pageSize).Execute() + savepointsPage, httpResponse, err = cmfClient.SavepointsApi.GetSavepointsForFlinkStatement(ctx, environment, statement).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() } else { - savepointsPage, httpResponse, err = cmfClient.SavepointsApi.GetSavepointsForFlinkApplication(ctx, environment, application).Page(currentPageNumber).Size(pageSize).Execute() + savepointsPage, httpResponse, err = cmfClient.SavepointsApi.GetSavepointsForFlinkApplication(ctx, environment, application).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() } if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list savepoints in the environment "%s": %s`, environment, parsedErr) @@ -412,7 +419,7 @@ func (cmfClient *CmfRestClient) ListDetachedSavepoint(ctx context.Context, filte var currentPageNumber int32 = 0 for !done { - savepointsPage, httpResponse, err := cmfClient.DetachedSavepointsApi.ListDetachedSavepoints(ctx).Page(currentPageNumber).Size(pageSize).Name(filter).Execute() + savepointsPage, httpResponse, err := cmfClient.DetachedSavepointsApi.ListDetachedSavepoints(ctx).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Name(filter).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list detached savepoints %s`, parsedErr) } @@ -461,7 +468,7 @@ func (cmfClient *CmfRestClient) ListComputePools(ctx context.Context, environmen var currentPageNumber int32 = 0 for !done { - computePoolsPage, httpResponse, err := cmfClient.SQLApi.GetComputePools(ctx, environment).Page(currentPageNumber).Size(pageSize).Execute() + computePoolsPage, httpResponse, err := cmfClient.SQLApi.GetComputePools(ctx, environment).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list compute pools in the environment "%s": %s`, environment, parsedErr) } @@ -509,7 +516,7 @@ func (cmfClient *CmfRestClient) ListStatements(ctx context.Context, environment, const pageSize = 100 var currentPageNumber int32 = 0 - request := cmfClient.SQLApi.GetStatements(ctx, environment) + request := cmfClient.SQLApi.GetStatements(ctx, environment).Sort(sortByName) if computePool != "" { request = request.ComputePool(computePool) } @@ -614,7 +621,7 @@ func (cmfClient *CmfRestClient) ListCatalog(ctx context.Context) ([]cmfsdk.Kafka var currentPageNumber int32 = 0 for !done { - catalogPage, httpResponse, err := cmfClient.SQLApi.GetKafkaCatalogs(ctx).Page(currentPageNumber).Size(pageSize).Execute() + catalogPage, httpResponse, err := cmfClient.SQLApi.GetKafkaCatalogs(ctx).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list Kafka Catalog: %s`, parsedErr) } @@ -654,7 +661,7 @@ func (cmfClient *CmfRestClient) ListApplicationInstances(ctx context.Context, en done := false for !done { - instancesPage, httpResponse, err := cmfClient.FlinkApplicationsApi.GetApplicationInstances(ctx, environment, application).Page(currentPageNumber).Size(pageSize).Execute() + instancesPage, httpResponse, err := cmfClient.FlinkApplicationsApi.GetApplicationInstances(ctx, environment, application).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list instances of application "%s" in the environment "%s": %s`, application, environment, parsedErr) } @@ -692,7 +699,8 @@ func (cmfClient *CmfRestClient) ListSecretMappings(ctx context.Context, envName var currentPageNumber int32 = 0 for !done { - mappingsPage, httpResponse, err := cmfClient.EnvironmentsApi.GetEnvironmentSecretMappings(ctx, envName).Page(currentPageNumber).Size(pageSize).Execute() + // "name" is not unique for mappings, so add "uid" (the primary key) as a tiebreaker. + mappingsPage, httpResponse, err := cmfClient.EnvironmentsApi.GetEnvironmentSecretMappings(ctx, envName).Sort([]string{"name", "uid"}).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list secret mappings in the environment "%s": %s`, envName, parsedErr) } @@ -740,7 +748,7 @@ func (cmfClient *CmfRestClient) ListSecrets(ctx context.Context) ([]cmfsdk.Secre var currentPageNumber int32 = 0 for !done { - secretsPage, httpResponse, err := cmfClient.SecretsApi.GetSecrets(ctx).Page(currentPageNumber).Size(pageSize).Execute() + secretsPage, httpResponse, err := cmfClient.SecretsApi.GetSecrets(ctx).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list secrets: %s`, parsedErr) } @@ -804,7 +812,7 @@ func (cmfClient *CmfRestClient) ListDatabases(ctx context.Context, catalogName s var currentPageNumber int32 = 0 for !done { - databasePage, httpResponse, err := cmfClient.SQLApi.GetKafkaDatabases(ctx, catalogName).Page(currentPageNumber).Size(pageSize).Execute() + databasePage, httpResponse, err := cmfClient.SQLApi.GetKafkaDatabases(ctx, catalogName).Sort(sortByName).Page(currentPageNumber).Size(pageSize).Execute() if parsedErr := parseSdkError(httpResponse, err); parsedErr != nil { return nil, fmt.Errorf(`failed to list databases in catalog "%s": %s`, catalogName, parsedErr) } diff --git a/pkg/flink/cmf_rest_client_test.go b/pkg/flink/cmf_rest_client_test.go new file mode 100644 index 0000000000..668adf1c2a --- /dev/null +++ b/pkg/flink/cmf_rest_client_test.go @@ -0,0 +1,136 @@ +package flink + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strconv" + "sync" + "testing" + + "github.com/stretchr/testify/require" + + cmfsdk "github.com/confluentinc/cmf-sdk-go/v1" +) + +// newTestCmfClient builds a CmfRestClient that talks to srv. +func newTestCmfClient(srv *httptest.Server) *CmfRestClient { + cfg := cmfsdk.NewConfiguration() + cfg.Servers = cmfsdk.ServerConfigurations{{URL: srv.URL}} + cfg.HTTPClient = srv.Client() + return &CmfRestClient{APIClient: cmfsdk.NewAPIClient(cfg)} +} + +// sortRecorder captures the sort query params of the most recent request. +type sortRecorder struct { + mu sync.Mutex + last []string +} + +func (s *sortRecorder) set(v []string) { + s.mu.Lock() + s.last = v + s.mu.Unlock() +} + +func (s *sortRecorder) get() []string { + s.mu.Lock() + defer s.mu.Unlock() + return s.last +} + +// TestListStatements_ReturnsAllItemsAcrossPages guards the paging regression: a +// paginated list of >100 items must return every item exactly once. The mock pages +// a stable order only when the request carries the unique "name" sort; without it +// the order rotates per page (as Postgres would with no ORDER BY), skipping and +// duplicating rows across page boundaries. +func TestListStatements_ReturnsAllItemsAcrossPages(t *testing.T) { + const total = 250 // spans page boundaries at 100 and 200 + + names := make([]string, total) + for i := range names { + names[i] = fmt.Sprintf("stmt-%04d", i) + } + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + q := r.URL.Query() + sortParam := q["sort"] + page, _ := strconv.Atoi(q.Get("page")) + size, _ := strconv.Atoi(q.Get("size")) + + // A unique sort makes paging deterministic. Without it, rotate each page's + // window so pages overlap and leave gaps, as unordered OFFSET paging would. + stable := len(sortParam) == 1 && sortParam[0] == "name" + + var items []cmfsdk.Statement + for i := page * size; i < min((page+1)*size, total); i++ { + idx := i + if !stable { + idx = (i + page*size) % total + } + items = append(items, cmfsdk.Statement{Metadata: cmfsdk.StatementMetadata{Name: names[idx]}}) + } + out := cmfsdk.StatementsPage{} + out.SetItems(items) + w.Header().Set("Content-Type", "application/json") + require.NoError(t, json.NewEncoder(w).Encode(out)) + })) + defer srv.Close() + + got, err := newTestCmfClient(srv).ListStatements(context.Background(), "env", "", "") + require.NoError(t, err) + + seen := make(map[string]int, total) + for _, s := range got { + seen[s.Metadata.Name]++ + } + require.Len(t, seen, total, "every statement returned exactly once") + for name, count := range seen { + require.Equal(t, 1, count, "statement %q returned %d times", name, count) + } +} + +// TestListCommands_SendUniqueSortKey locks in the per-resource sort key sent by +// every paginated list call. An empty page ends the loop after one request, so +// the recorder holds exactly that call's sort params. +func TestListCommands_SendUniqueSortKey(t *testing.T) { + rec := &sortRecorder{} + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + rec.set(r.URL.Query()["sort"]) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"items":[]}`)) + })) + defer srv.Close() + + c := newTestCmfClient(srv) + ctx := context.Background() + + cases := []struct { + name string + invoke func() error + want []string + }{ + {"statements", func() error { _, e := c.ListStatements(ctx, "env", "", ""); return e }, []string{"name"}}, + {"applications", func() error { _, e := c.ListApplications(ctx, "env"); return e }, []string{"name"}}, + {"application-instances", func() error { _, e := c.ListApplicationInstances(ctx, "env", "app"); return e }, []string{"name"}}, + {"application-events", func() error { _, e := c.ListApplicationEvents(ctx, "env", "app"); return e }, []string{"creationTimestamp,desc", "name"}}, + {"savepoints-statement", func() error { _, e := c.ListSavepoint(ctx, "env", "stmt", "", true); return e }, []string{"name"}}, + {"savepoints-application", func() error { _, e := c.ListSavepoint(ctx, "env", "", "app", false); return e }, []string{"name"}}, + {"detached-savepoints", func() error { _, e := c.ListDetachedSavepoint(ctx, ""); return e }, []string{"name"}}, + {"compute-pools", func() error { _, e := c.ListComputePools(ctx, "env"); return e }, []string{"name"}}, + {"environments", func() error { _, e := c.ListEnvironments(ctx); return e }, []string{"name"}}, + {"catalogs", func() error { _, e := c.ListCatalog(ctx); return e }, []string{"name"}}, + {"databases", func() error { _, e := c.ListDatabases(ctx, "cat"); return e }, []string{"name"}}, + {"secrets", func() error { _, e := c.ListSecrets(ctx); return e }, []string{"name"}}, + {"secret-mappings", func() error { _, e := c.ListSecretMappings(ctx, "env"); return e }, []string{"name", "uid"}}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + rec.set(nil) + require.NoError(t, tc.invoke()) + require.Equal(t, tc.want, rec.get()) + }) + } +}