[api] Add ListKubernetesClusters (#59820)

This commits deperecates GetKubernetesClusters, which
should be replaced by a paginated counterpart:
ListKubernetesClusters

Towards: gravitational/teleport.e/issues/6759
Changelog: Add paginated API ListKubernetesClusters, deprecate GetKubernetesClusters
This commit is contained in:
Luke Okraszewski
2025-10-08 15:01:19 +00:00
committed by GitHub
parent 89a13c7c7c
commit fca48f3b9b
18 changed files with 1859 additions and 1003 deletions
+23
View File
@@ -3479,7 +3479,9 @@ func (c *Client) GetKubernetesCluster(ctx context.Context, name string) (types.K
}
// GetKubernetesClusters returns all kubernetes cluster resources.
// Deprecated: Prefer paginated variant such as [ListKubernetesClusters] or [RangeKubernetesClusters]
func (c *Client) GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster, error) {
//nolint:staticcheck // TODO(okraport): deprecated, to be removed in v21
items, err := c.grpc.GetKubernetesClusters(ctx, &emptypb.Empty{})
if err != nil {
return nil, trace.Wrap(err)
@@ -3491,6 +3493,27 @@ func (c *Client) GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster
return clusters, nil
}
// ListKubernetesClusters returns a page of registered kubernetes clusters.
func (c *Client) ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error) {
resp, err := c.grpc.ListKubernetesClusters(ctx, &proto.ListKubernetesClustersRequest{
PageSize: int32(limit),
PageToken: start,
})
if err != nil {
return nil, "", trace.Wrap(err)
}
kubeClusters := make([]types.KubeCluster, len(resp.KubernetesClusters))
for i := range resp.KubernetesClusters {
kubeClusters[i] = resp.KubernetesClusters[i]
}
return kubeClusters, resp.NextPageToken, nil
}
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
func (c *Client) RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error] {
return clientutils.RangeResources(ctx, start, end, c.ListKubernetesClusters, types.KubeCluster.GetName)
}
// DeleteKubernetesCluster deletes specified kubernetes cluster resource.
func (c *Client) DeleteKubernetesCluster(ctx context.Context, name string) error {
_, err := c.grpc.DeleteKubernetesCluster(ctx, &types.ResourceRequest{Name: name})
File diff suppressed because it is too large Load Diff
+43
View File
@@ -235,6 +235,7 @@ const (
AuthService_DeleteDatabase_FullMethodName = "/proto.AuthService/DeleteDatabase"
AuthService_DeleteAllDatabases_FullMethodName = "/proto.AuthService/DeleteAllDatabases"
AuthService_GetKubernetesClusters_FullMethodName = "/proto.AuthService/GetKubernetesClusters"
AuthService_ListKubernetesClusters_FullMethodName = "/proto.AuthService/ListKubernetesClusters"
AuthService_GetKubernetesCluster_FullMethodName = "/proto.AuthService/GetKubernetesCluster"
AuthService_CreateKubernetesCluster_FullMethodName = "/proto.AuthService/CreateKubernetesCluster"
AuthService_UpdateKubernetesCluster_FullMethodName = "/proto.AuthService/UpdateKubernetesCluster"
@@ -830,8 +831,11 @@ type AuthServiceClient interface {
DeleteDatabase(ctx context.Context, in *types.ResourceRequest, opts ...grpc.CallOption) (*emptypb.Empty, error)
// DeleteAllDatabases removes all database resources.
DeleteAllDatabases(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*emptypb.Empty, error)
// Deprecated: Do not use.
// GetKubernetesClusters returns all registered kubernetes clusters.
GetKubernetesClusters(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*types.KubernetesClusterV3List, error)
// ListKubernetesClusters returns a page of registered kubernetes clusters.
ListKubernetesClusters(ctx context.Context, in *ListKubernetesClustersRequest, opts ...grpc.CallOption) (*ListKubernetesClustersResponse, error)
// GetKubernetesCluster returns a kubernetes cluster by name.
GetKubernetesCluster(ctx context.Context, in *types.ResourceRequest, opts ...grpc.CallOption) (*types.KubernetesClusterV3, error)
// CreateKubernetesCluster creates a new kubernetes cluster resource.
@@ -3165,6 +3169,7 @@ func (c *authServiceClient) DeleteAllDatabases(ctx context.Context, in *emptypb.
return out, nil
}
// Deprecated: Do not use.
func (c *authServiceClient) GetKubernetesClusters(ctx context.Context, in *emptypb.Empty, opts ...grpc.CallOption) (*types.KubernetesClusterV3List, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(types.KubernetesClusterV3List)
@@ -3175,6 +3180,16 @@ func (c *authServiceClient) GetKubernetesClusters(ctx context.Context, in *empty
return out, nil
}
func (c *authServiceClient) ListKubernetesClusters(ctx context.Context, in *ListKubernetesClustersRequest, opts ...grpc.CallOption) (*ListKubernetesClustersResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(ListKubernetesClustersResponse)
err := c.cc.Invoke(ctx, AuthService_ListKubernetesClusters_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
func (c *authServiceClient) GetKubernetesCluster(ctx context.Context, in *types.ResourceRequest, opts ...grpc.CallOption) (*types.KubernetesClusterV3, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(types.KubernetesClusterV3)
@@ -4380,8 +4395,11 @@ type AuthServiceServer interface {
DeleteDatabase(context.Context, *types.ResourceRequest) (*emptypb.Empty, error)
// DeleteAllDatabases removes all database resources.
DeleteAllDatabases(context.Context, *emptypb.Empty) (*emptypb.Empty, error)
// Deprecated: Do not use.
// GetKubernetesClusters returns all registered kubernetes clusters.
GetKubernetesClusters(context.Context, *emptypb.Empty) (*types.KubernetesClusterV3List, error)
// ListKubernetesClusters returns a page of registered kubernetes clusters.
ListKubernetesClusters(context.Context, *ListKubernetesClustersRequest) (*ListKubernetesClustersResponse, error)
// GetKubernetesCluster returns a kubernetes cluster by name.
GetKubernetesCluster(context.Context, *types.ResourceRequest) (*types.KubernetesClusterV3, error)
// CreateKubernetesCluster creates a new kubernetes cluster resource.
@@ -5186,6 +5204,9 @@ func (UnimplementedAuthServiceServer) DeleteAllDatabases(context.Context, *empty
func (UnimplementedAuthServiceServer) GetKubernetesClusters(context.Context, *emptypb.Empty) (*types.KubernetesClusterV3List, error) {
return nil, status.Errorf(codes.Unimplemented, "method GetKubernetesClusters not implemented")
}
func (UnimplementedAuthServiceServer) ListKubernetesClusters(context.Context, *ListKubernetesClustersRequest) (*ListKubernetesClustersResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method ListKubernetesClusters not implemented")
}
func (UnimplementedAuthServiceServer) GetKubernetesCluster(context.Context, *types.ResourceRequest) (*types.KubernetesClusterV3, error) {
return nil, status.Errorf(codes.Unimplemented, "method GetKubernetesCluster not implemented")
}
@@ -8842,6 +8863,24 @@ func _AuthService_GetKubernetesClusters_Handler(srv interface{}, ctx context.Con
return interceptor(ctx, in, info, handler)
}
func _AuthService_ListKubernetesClusters_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(ListKubernetesClustersRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(AuthServiceServer).ListKubernetesClusters(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: AuthService_ListKubernetesClusters_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(AuthServiceServer).ListKubernetesClusters(ctx, req.(*ListKubernetesClustersRequest))
}
return interceptor(ctx, in, info, handler)
}
func _AuthService_GetKubernetesCluster_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(types.ResourceRequest)
if err := dec(in); err != nil {
@@ -10776,6 +10815,10 @@ var AuthService_ServiceDesc = grpc.ServiceDesc{
MethodName: "GetKubernetesClusters",
Handler: _AuthService_GetKubernetesClusters_Handler,
},
{
MethodName: "ListKubernetesClusters",
Handler: _AuthService_ListKubernetesClusters_Handler,
},
{
MethodName: "GetKubernetesCluster",
Handler: _AuthService_GetKubernetesCluster_Handler,
@@ -2781,6 +2781,22 @@ message ValidateTrustedClusterResponse {
repeated types.CertAuthorityV2 CertAuthorities = 1;
}
message ListKubernetesClustersRequest {
// The maximum number of items to return.
// The server may impose a different page size at its discretion.
int32 page_size = 1;
// The next_page_token value returned from a previous List request, if any.
string page_token = 2;
}
message ListKubernetesClustersResponse {
// KubernetesClusters is a list of kubernetes cluster resources.
repeated types.KubernetesClusterV3 kubernetes_clusters = 1;
// Token to retrieve the next page of results, or empty if there are no
// more results in the list.
string next_page_token = 2;
}
// AuthService is authentication/authorization service implementation
service AuthService {
// InventoryControlStream is the per-instance stream used to advertise teleport instance
@@ -3400,7 +3416,11 @@ service AuthService {
rpc DeleteAllDatabases(google.protobuf.Empty) returns (google.protobuf.Empty);
// GetKubernetesClusters returns all registered kubernetes clusters.
rpc GetKubernetesClusters(google.protobuf.Empty) returns (types.KubernetesClusterV3List);
rpc GetKubernetesClusters(google.protobuf.Empty) returns (types.KubernetesClusterV3List) {
option deprecated = true;
}
// ListKubernetesClusters returns a page of registered kubernetes clusters.
rpc ListKubernetesClusters(ListKubernetesClustersRequest) returns (ListKubernetesClustersResponse);
// GetKubernetesCluster returns a kubernetes cluster by name.
rpc GetKubernetesCluster(types.ResourceRequest) returns (types.KubernetesClusterV3);
// CreateKubernetesCluster creates a new kubernetes cluster resource.
@@ -132,7 +132,6 @@ func newDefaultConfig() *Config {
"proto.AuthService.GetEvents": {},
"proto.AuthService.GetGithubConnectors": {},
"proto.AuthService.GetInstallers": {},
"proto.AuthService.GetKubernetesClusters": {},
"proto.AuthService.GetLocks": {},
"proto.AuthService.GetMFADevices": {},
"proto.AuthService.GetOIDCConnectors": {},
+57 -7
View File
@@ -6492,17 +6492,67 @@ func (a *ServerWithRoles) GetKubernetesClusters(ctx context.Context) (result []t
if err := a.authorizeAction(types.KindKubernetesCluster, types.VerbList, types.VerbRead); err != nil {
return nil, trace.Wrap(err)
}
// Filter out kube clusters user doesn't have access to.
clusters, err := a.authServer.GetKubernetesClusters(ctx)
out, err := iterstream.Collect(
iterstream.FilterMap(
a.authServer.RangeKubernetesClusters(ctx, "", ""),
func(cluster types.KubeCluster) (types.KubeCluster, bool) {
// Filter out kube clusters user doesn't have access to.
if a.checkAccessToKubeCluster(cluster) == nil {
return cluster, true
}
return nil, false
},
),
)
if err != nil {
return nil, trace.Wrap(err)
}
for _, cluster := range clusters {
if err := a.checkAccessToKubeCluster(cluster); err == nil {
result = append(result, cluster)
}
return out, nil
}
// ListKubernetesClusters returns a page of registered kubernetes clusters.
func (a *ServerWithRoles) ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error) {
if err := a.authorizeAction(types.KindKubernetesCluster, types.VerbList, types.VerbRead); err != nil {
return nil, "", trace.Wrap(err)
}
return result, nil
if limit <= 0 || limit > apidefaults.DefaultChunkSize {
limit = apidefaults.DefaultChunkSize
}
var next string
var seen int
out, err := iterstream.Collect(
iterstream.TakeWhile(
iterstream.FilterMap(
a.authServer.RangeKubernetesClusters(ctx, start, ""),
func(cluster types.KubeCluster) (types.KubeCluster, bool) {
// Filter out kube clusters user doesn't have access to.
if a.checkAccessToKubeCluster(cluster) == nil {
return cluster, true
}
return nil, false
},
),
func(cluster types.KubeCluster) bool {
if seen < limit {
seen++
return true
}
next = cluster.GetName()
return false
},
),
)
if err != nil {
return nil, "", trace.Wrap(err)
}
return out, next, nil
}
// DeleteKubernetesCluster removes the specified kubernetes cluster resource.
+21 -1
View File
@@ -3045,10 +3045,22 @@ func TestKubernetesClusterCRUD_DiscoveryService(t *testing.T) {
require.True(t, trace.IsAccessDenied(discoveryClt.CreateKubernetesCluster(ctx, clusterWithDynamicLabels)))
})
t.Run("Read", func(t *testing.T) {
diffopt := cmpopts.IgnoreFields(types.Metadata{}, "Revision")
clusters, err := discoveryClt.GetKubernetesClusters(ctx)
require.NoError(t, err)
require.Empty(t, cmp.Diff([]types.KubeCluster{eksCluster}, clusters, cmpopts.IgnoreFields(types.Metadata{}, "Revision")))
require.Empty(t, cmp.Diff([]types.KubeCluster{eksCluster}, clusters, diffopt))
clusters, next, err := discoveryClt.ListKubernetesClusters(ctx, 0, "")
require.Empty(t, next)
require.NoError(t, err)
require.Empty(t, cmp.Diff([]types.KubeCluster{eksCluster}, clusters, diffopt))
clusters, err = stream.Collect(discoveryClt.RangeKubernetesClusters(ctx, "", ""))
require.NoError(t, err)
require.Empty(t, cmp.Diff([]types.KubeCluster{eksCluster}, clusters, diffopt))
})
t.Run("Update", func(t *testing.T) {
require.NoError(t, discoveryClt.UpdateKubernetesCluster(ctx, eksCluster))
require.True(t, trace.IsAccessDenied(discoveryClt.UpdateKubernetesCluster(ctx, nonCloudCluster)))
@@ -3059,10 +3071,18 @@ func TestKubernetesClusterCRUD_DiscoveryService(t *testing.T) {
require.NoError(t, err)
require.Empty(t, clusters)
clusters, err = stream.Collect(discoveryClt.RangeKubernetesClusters(ctx, "", ""))
require.NoError(t, err)
require.Empty(t, clusters)
// Discovery service cannot delete non-cloud clusters.
clusters, err = srv.Auth().GetKubernetesClusters(ctx)
require.NoError(t, err)
require.Len(t, clusters, 1)
clusters, err = stream.Collect(srv.Auth().RangeKubernetesClusters(ctx, "", ""))
require.NoError(t, err)
require.Len(t, clusters, 1)
})
}
+16
View File
@@ -306,6 +306,10 @@ type ReadProxyAccessPoint interface {
// GetKubernetesClusters returns all kubernetes cluster resources.
GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster, error)
// ListKubernetesClusters returns a page of registered kubernetes clusters.
ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error)
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error]
// GetKubernetesCluster returns the specified kubernetes cluster resource.
GetKubernetesCluster(ctx context.Context, name string) (types.KubeCluster, error)
@@ -535,6 +539,10 @@ type ReadKubernetesAccessPoint interface {
// GetKubernetesClusters returns all kubernetes cluster resources.
GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster, error)
// ListKubernetesClusters returns a page of registered kubernetes clusters.
ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error)
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error]
// GetKubernetesCluster returns the specified kubernetes cluster resource.
GetKubernetesCluster(ctx context.Context, name string) (types.KubeCluster, error)
}
@@ -781,6 +789,10 @@ type ReadDiscoveryAccessPoint interface {
GetKubernetesCluster(ctx context.Context, name string) (types.KubeCluster, error)
// GetKubernetesClusters returns all kubernetes cluster resources.
GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster, error)
// ListKubernetesClusters returns a page of registered kubernetes clusters.
ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error)
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error]
// GetKubernetesServers returns all registered kubernetes servers.
GetKubernetesServers(ctx context.Context) ([]types.KubeServer, error)
@@ -1209,6 +1221,10 @@ type Cache interface {
// GetKubernetesClusters returns all kubernetes cluster resources.
GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster, error)
// ListKubernetesClusters returns a page of registered kubernetes clusters.
ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error)
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error]
// GetKubernetesCluster returns the specified kubernetes cluster resource.
GetKubernetesCluster(ctx context.Context, name string) (types.KubeCluster, error)
+28
View File
@@ -5039,6 +5039,34 @@ func (g *GRPCServer) GetKubernetesClusters(ctx context.Context, _ *emptypb.Empty
}, nil
}
// ListKubernetesClusters returns a page of registered kubernetes clusters.
func (g *GRPCServer) ListKubernetesClusters(ctx context.Context, req *authpb.ListKubernetesClustersRequest) (*authpb.ListKubernetesClustersResponse, error) {
auth, err := g.authenticate(ctx)
if err != nil {
return nil, trace.Wrap(err)
}
clusters, next, err := auth.ListKubernetesClusters(ctx, int(req.PageSize), req.PageToken)
if err != nil {
return nil, trace.Wrap(err)
}
resp := &authpb.ListKubernetesClustersResponse{
KubernetesClusters: make([]*types.KubernetesClusterV3, 0, len(clusters)),
NextPageToken: next,
}
for _, cluster := range clusters {
clusterV3, ok := cluster.(*types.KubernetesClusterV3)
if !ok {
return nil, trace.BadParameter("unsupported kubernetes cluster type %T", clusterV3)
}
resp.KubernetesClusters = append(resp.KubernetesClusters, clusterV3)
}
return resp, nil
}
// DeleteKubernetesCluster removes the specified kubernetes cluster.
func (g *GRPCServer) DeleteKubernetesCluster(ctx context.Context, req *types.ResourceRequest) (*emptypb.Empty, error) {
auth, err := g.authenticate(ctx)
+49
View File
@@ -18,6 +18,7 @@ package cache
import (
"context"
"iter"
"github.com/gravitational/trace"
"google.golang.org/protobuf/proto"
@@ -147,6 +148,54 @@ func (c *Cache) GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster,
return out, nil
}
// ListKubernetesClusters returns a page of registered kubernetes clusters.
func (c *Cache) ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error) {
ctx, span := c.Tracer.Start(ctx, "cache/ListKubernetesClusters")
defer span.End()
lister := genericLister[types.KubeCluster, kubeClusterIndex]{
cache: c,
collection: c.collections.kubeClusters,
index: kubeClusterNameIndex,
upstreamList: c.Config.Kubernetes.ListKubernetesClusters,
nextToken: types.KubeCluster.GetName,
}
out, next, err := lister.list(ctx, limit, start)
if err != nil {
return nil, "", trace.Wrap(err)
}
return out, next, nil
}
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
func (c *Cache) RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error] {
lister := genericLister[types.KubeCluster, kubeClusterIndex]{
cache: c,
collection: c.collections.kubeClusters,
index: kubeClusterNameIndex,
upstreamList: c.Config.Kubernetes.ListKubernetesClusters,
nextToken: types.KubeCluster.GetName,
// TODO(lokraszewski): DELETE IN v21.0.0
fallbackGetter: c.Config.Kubernetes.GetKubernetesClusters,
}
return func(yield func(types.KubeCluster, error) bool) {
ctx, span := c.Tracer.Start(ctx, "cache/RangeKubernetesClusters")
defer span.End()
for cluster, err := range lister.RangeWithFallback(ctx, start, end) {
if !yield(cluster, err) {
return
}
if err != nil {
return
}
}
}
}
// GetKubernetesCluster returns the specified kubernetes cluster resource.
func (c *Cache) GetKubernetesCluster(ctx context.Context, name string) (types.KubeCluster, error) {
ctx, span := c.Tracer.Start(ctx, "cache/GetKubernetesCluster")
+33 -13
View File
@@ -35,19 +35,39 @@ func TestKubernetes(t *testing.T) {
p := newTestPack(t, ForProxy)
t.Cleanup(p.Close)
testResources(t, p, testFuncs[types.KubeCluster]{
newResource: func(name string) (types.KubeCluster, error) {
return types.NewKubernetesClusterV3(types.Metadata{
Name: name,
}, types.KubernetesClusterSpecV3{})
},
create: p.kubernetes.CreateKubernetesCluster,
list: getAllAdapter(p.kubernetes.GetKubernetesClusters),
cacheGet: p.cache.GetKubernetesCluster,
cacheList: getAllAdapter(p.cache.GetKubernetesClusters),
update: p.kubernetes.UpdateKubernetesCluster,
deleteAll: p.kubernetes.DeleteAllKubernetesClusters,
}, withSkipPaginationTest())
t.Run("GetKubernetesClusters", func(t *testing.T) {
testResources(t, p, testFuncs[types.KubeCluster]{
newResource: func(name string) (types.KubeCluster, error) {
return types.NewKubernetesClusterV3(types.Metadata{
Name: name,
}, types.KubernetesClusterSpecV3{})
},
create: p.kubernetes.CreateKubernetesCluster,
list: getAllAdapter(p.kubernetes.GetKubernetesClusters),
cacheGet: p.cache.GetKubernetesCluster,
cacheList: getAllAdapter(p.cache.GetKubernetesClusters),
update: p.kubernetes.UpdateKubernetesCluster,
deleteAll: p.kubernetes.DeleteAllKubernetesClusters,
}, withSkipPaginationTest())
})
t.Run("ListKubernetesClusters", func(t *testing.T) {
testResources(t, p, testFuncs[types.KubeCluster]{
newResource: func(name string) (types.KubeCluster, error) {
return types.NewKubernetesClusterV3(types.Metadata{
Name: name,
}, types.KubernetesClusterSpecV3{})
},
create: p.kubernetes.CreateKubernetesCluster,
list: p.kubernetes.ListKubernetesClusters,
cacheGet: p.cache.GetKubernetesCluster,
cacheList: p.cache.ListKubernetesClusters,
update: p.kubernetes.UpdateKubernetesCluster,
deleteAll: p.kubernetes.DeleteAllKubernetesClusters,
Range: p.kubernetes.RangeKubernetesClusters,
cacheRange: p.cache.RangeKubernetesClusters,
})
})
}
// TestKubernetesServers tests that CRUD operations on kube servers are
+5
View File
@@ -20,6 +20,7 @@ package services
import (
"context"
"iter"
"github.com/gravitational/trace"
@@ -31,6 +32,10 @@ import (
type KubernetesClusterGetter interface {
// GetKubernetesClusters returns all kubernetes cluster resources.
GetKubernetesClusters(context.Context) ([]types.KubeCluster, error)
// ListKubernetesClusters returns a page of registered kubernetes clusters.
ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error)
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error]
// GetKubernetesCluster returns the specified kubernetes cluster resource.
GetKubernetesCluster(ctx context.Context, name string) (types.KubeCluster, error)
}
+81 -12
View File
@@ -20,41 +20,110 @@ package local
import (
"context"
"iter"
"log/slog"
"github.com/gravitational/trace"
"github.com/gravitational/teleport"
"github.com/gravitational/teleport/api/defaults"
"github.com/gravitational/teleport/api/types"
"github.com/gravitational/teleport/lib/backend"
"github.com/gravitational/teleport/lib/itertools/stream"
"github.com/gravitational/teleport/lib/services"
)
// KubernetesService manages kubernetes resources in the backend.
type KubernetesService struct {
backend.Backend
logger *slog.Logger
}
// NewKubernetesService creates a new KubernetesService.
func NewKubernetesService(backend backend.Backend) *KubernetesService {
return &KubernetesService{Backend: backend}
return &KubernetesService{
Backend: backend,
logger: slog.With(teleport.ComponentKey, "KubernetesService"),
}
}
// GetKubernetesClusters returns all kubernetes cluster resources.
func (s *KubernetesService) GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster, error) {
startKey := backend.ExactKey(kubernetesPrefix)
result, err := s.GetRange(ctx, startKey, backend.RangeEnd(startKey), backend.NoLimit)
out, err := stream.Collect(s.RangeKubernetesClusters(ctx, "", ""))
if err != nil {
return nil, trace.Wrap(err)
}
kubeClusters := make([]types.KubeCluster, len(result.Items))
for i, item := range result.Items {
cluster, err := services.UnmarshalKubeCluster(item.Value,
services.WithExpires(item.Expires), services.WithRevision(item.Revision))
if err != nil {
return nil, trace.Wrap(err)
}
kubeClusters[i] = cluster
return out, nil
}
// ListKubernetesClusters returns a page of registered kubernetes clusters.
func (s *KubernetesService) ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error) {
// Adjust page size, so it can't be too large.
if limit <= 0 || limit > defaults.DefaultChunkSize {
limit = defaults.DefaultChunkSize
}
return kubeClusters, nil
var next string
var seen int
out, err := stream.Collect(
stream.TakeWhile(
s.RangeKubernetesClusters(ctx, start, ""),
func(cluster types.KubeCluster) bool {
if seen < limit {
seen++
return true
}
next = cluster.GetName()
return false
},
),
)
if err != nil {
return nil, "", trace.Wrap(err)
}
return out, next, nil
}
// RangeKubernetesClusters returns kubernetes clusters within the range [start, end).
func (s *KubernetesService) RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error] {
mapFn := func(item backend.Item) (types.KubeCluster, bool) {
cluster, err := services.UnmarshalKubeCluster(item.Value,
services.WithExpires(item.Expires),
services.WithRevision(item.Revision))
if err != nil {
s.logger.WarnContext(ctx, "Failed to unmarshal kubernetes cluster",
"key", item.Key,
"error", err,
)
return nil, false
}
return cluster, true
}
kubernetesKey := backend.NewKey(kubernetesPrefix)
startKey := kubernetesKey.AppendKey(backend.KeyFromString(start))
endKey := backend.RangeEnd(kubernetesKey)
if end != "" {
endKey = kubernetesKey.AppendKey(backend.KeyFromString(end)).ExactKey()
}
return stream.TakeWhile(
stream.FilterMap(
s.Backend.Items(ctx, backend.ItemsParams{
StartKey: startKey,
EndKey: endKey,
}),
mapFn,
),
func(cluster types.KubeCluster) bool {
// The range is not inclusive of the end key, so return early
// if the end has been reached.
return end == "" || cluster.GetName() < end
})
}
// GetKubernetesCluster returns the specified kubernetes cluster resource.
+44 -12
View File
@@ -30,6 +30,7 @@ import (
"github.com/gravitational/teleport/api/types"
"github.com/gravitational/teleport/lib/backend/memory"
"github.com/gravitational/teleport/lib/itertools/stream"
)
// TestKubernetesCRUD tests backend operations with kubernetes resources.
@@ -53,6 +54,10 @@ func TestKubernetesCRUD(t *testing.T) {
Name: "c2",
}, types.KubernetesClusterSpecV3{})
require.NoError(t, err)
kubeCluster3, err := types.NewKubernetesClusterV3(types.Metadata{
Name: "c3",
}, types.KubernetesClusterSpecV3{})
require.NoError(t, err)
// Initially we expect no Kubernetess.
out, err := service.GetKubernetesClusters(ctx)
@@ -64,20 +69,49 @@ func TestKubernetesCRUD(t *testing.T) {
require.NoError(t, err)
err = service.CreateKubernetesCluster(ctx, kubeCluster2)
require.NoError(t, err)
err = service.CreateKubernetesCluster(ctx, kubeCluster3)
require.NoError(t, err)
expectedAll := []types.KubeCluster{kubeCluster1, kubeCluster2, kubeCluster3}
diffopt := cmpopts.IgnoreFields(types.Metadata{}, "Revision")
// Fetch all Kubernetess.
out, err = service.GetKubernetesClusters(ctx)
require.NoError(t, err)
require.Empty(t, cmp.Diff([]types.KubeCluster{kubeCluster1, kubeCluster2}, out,
cmpopts.IgnoreFields(types.Metadata{}, "Revision"),
))
require.Empty(t, cmp.Diff(expectedAll, out, diffopt))
// List with page limit
page1, page2Start, err := service.ListKubernetesClusters(ctx, 2, "")
require.NoError(t, err)
require.NotEmpty(t, page2Start)
require.Len(t, page1, 2)
// List with start
page2, next, err := service.ListKubernetesClusters(ctx, 2, page2Start)
require.NoError(t, err)
require.Empty(t, next)
require.Len(t, page2, 1)
require.Empty(t, cmp.Diff(expectedAll, append(page1, page2...), diffopt))
// Range over all
out, err = stream.Collect(service.RangeKubernetesClusters(ctx, "", ""))
require.NoError(t, err)
require.Empty(t, cmp.Diff(expectedAll, out, diffopt))
// Range with upper bound
out, err = stream.Collect(service.RangeKubernetesClusters(ctx, "", page2Start))
require.NoError(t, err)
require.Empty(t, cmp.Diff(page1, out, diffopt))
// Range with lower bound
out, err = stream.Collect(service.RangeKubernetesClusters(ctx, page2Start, ""))
require.NoError(t, err)
require.Empty(t, cmp.Diff(page2, out, diffopt))
// Fetch a specific Kubernetes.
cluster, err := service.GetKubernetesCluster(ctx, kubeCluster2.GetName())
require.NoError(t, err)
require.Empty(t, cmp.Diff(kubeCluster2, cluster,
cmpopts.IgnoreFields(types.Metadata{}, "Revision"),
))
require.Empty(t, cmp.Diff(kubeCluster2, cluster, diffopt))
// Try to fetch a Kubernetes that doesn't exist.
_, err = service.GetKubernetesCluster(ctx, "doesnotexist")
@@ -93,18 +127,16 @@ func TestKubernetesCRUD(t *testing.T) {
require.NoError(t, err)
cluster, err = service.GetKubernetesCluster(ctx, kubeCluster1.GetName())
require.NoError(t, err)
require.Empty(t, cmp.Diff(kubeCluster1, cluster,
cmpopts.IgnoreFields(types.Metadata{}, "Revision"),
))
require.Empty(t, cmp.Diff(kubeCluster1, cluster, diffopt))
// Delete a Kubernetes.
err = service.DeleteKubernetesCluster(ctx, kubeCluster1.GetName())
require.NoError(t, err)
expectedAll = []types.KubeCluster{kubeCluster2, kubeCluster3}
out, err = service.GetKubernetesClusters(ctx)
require.NoError(t, err)
require.Empty(t, cmp.Diff([]types.KubeCluster{kubeCluster2}, out,
cmpopts.IgnoreFields(types.Metadata{}, "Revision"),
))
require.Empty(t, cmp.Diff(expectedAll, out, diffopt))
// Try to delete a Kubernetes that doesn't exist.
err = service.DeleteKubernetesCluster(ctx, "doesnotexist")
+2 -1
View File
@@ -36,6 +36,7 @@ import (
apiutils "github.com/gravitational/teleport/api/utils"
"github.com/gravitational/teleport/api/utils/retryutils"
"github.com/gravitational/teleport/lib/defaults"
iterstream "github.com/gravitational/teleport/lib/itertools/stream"
"github.com/gravitational/teleport/lib/services/readonly"
"github.com/gravitational/teleport/lib/utils"
logutils "github.com/gravitational/teleport/lib/utils/log"
@@ -528,7 +529,7 @@ func NewKubeClusterWatcher(ctx context.Context, cfg KubeClusterWatcherConfig) (*
ResourceWatcherConfig: cfg.ResourceWatcherConfig,
ResourceKind: types.KindKubernetesCluster,
ResourceGetter: func(ctx context.Context) ([]types.KubeCluster, error) {
return cfg.KubernetesClusterGetter.GetKubernetesClusters(ctx)
return iterstream.Collect(cfg.KubernetesClusterGetter.RangeKubernetesClusters(ctx, "", ""))
},
ResourceKey: types.KubeCluster.GetName,
ResourcesC: cfg.KubeClustersC,
+9
View File
@@ -347,6 +347,15 @@ func (m *mockAuthServer) UpdateDiscoveryConfigStatus(ctx context.Context, name s
func (m *mockAuthServer) GetKubernetesClusters(ctx context.Context) ([]types.KubeCluster, error) {
return nil, nil
}
func (m *mockAuthServer) ListKubernetesClusters(ctx context.Context, limit int, start string) ([]types.KubeCluster, string, error) {
return nil, "", nil
}
func (m *mockAuthServer) RangeKubernetesClusters(ctx context.Context, start, end string) iter.Seq2[types.KubeCluster, error] {
return func(yield func(types.KubeCluster, error) bool) {}
}
func (m *mockAuthServer) GetKubernetesServers(context.Context) ([]types.KubeServer, error) {
return nil, nil
}
@@ -39,6 +39,7 @@ import (
"github.com/gravitational/teleport/api/types"
"github.com/gravitational/teleport/lib/automaticupgrades"
"github.com/gravitational/teleport/lib/automaticupgrades/version"
iterstream "github.com/gravitational/teleport/lib/itertools/stream"
kubeutils "github.com/gravitational/teleport/lib/kube/utils"
"github.com/gravitational/teleport/lib/srv/discovery/common"
libslices "github.com/gravitational/teleport/lib/utils/slices"
@@ -127,7 +128,7 @@ func (s *Server) startKubeIntegrationWatchers() error {
continue
}
existingClusters, err := clt.GetKubernetesClusters(s.ctx)
existingClusters, err := iterstream.Collect(clt.RangeKubernetesClusters(s.ctx, "", ""))
if err != nil {
s.Log.WarnContext(s.ctx, "Failed to get Kubernetes clusters from cache", "error", err)
continue
+4 -2
View File
@@ -2036,7 +2036,8 @@ func (rc *ResourceCommand) Delete(ctx context.Context, client *authclient.Client
}
fmt.Printf("%s %q has been deleted\n", resDesc, name)
case types.KindKubernetesCluster:
clusters, err := client.GetKubernetesClusters(ctx)
// TODO(okraport) DELETE IN v21.0.0, replace with regular Collect
clusters, err := clientutils.CollectWithFallback(ctx, client.ListKubernetesClusters, client.GetKubernetesClusters)
if err != nil {
return trace.Wrap(err)
}
@@ -2865,7 +2866,8 @@ func (rc *ResourceCommand) getCollection(ctx context.Context, client *authclient
}
return &databaseCollection{databases: databases}, nil
case types.KindKubernetesCluster:
clusters, err := client.GetKubernetesClusters(ctx)
// TODO(okraport) DELETE IN v21.0.0, replace with regular Collect
clusters, err := clientutils.CollectWithFallback(ctx, client.ListKubernetesClusters, client.GetKubernetesClusters)
if err != nil {
return nil, trace.Wrap(err)
}