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
34 changes: 21 additions & 13 deletions pkg/flink/cmf_rest_client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Expand Down
136 changes: 136 additions & 0 deletions pkg/flink/cmf_rest_client_test.go
Original file line number Diff line number Diff line change
@@ -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())
})
}
}