mirror of
https://github.com/gravitational/teleport.git
synced 2026-09-24 16:17:11 +08:00
Remove deprecated node http endpoints (#12970)
This commit is contained in:
@@ -1707,17 +1707,6 @@ func (c *Client) GetNodes(ctx context.Context, namespace string) ([]types.Server
|
||||
Namespace: namespace,
|
||||
})
|
||||
if err != nil {
|
||||
// Underlying ListResources for nodes was not available, use fallback.
|
||||
//
|
||||
// DELETE IN 11.0.0
|
||||
if trace.IsNotImplemented(err) {
|
||||
servers, err := GetNodesWithLabels(ctx, c, namespace, nil)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
return servers, nil
|
||||
}
|
||||
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
@@ -1729,116 +1718,6 @@ func (c *Client) GetNodes(ctx context.Context, namespace string) ([]types.Server
|
||||
return servers, nil
|
||||
}
|
||||
|
||||
// NodeClient is an interface used by GetNodesWithLabels to abstract over implementations of
|
||||
// the ListNodes method.
|
||||
//
|
||||
// DELETE IN 11.0.0 with GetNodesWithLabels (used in both api/client/client.go and lib/auth/httpfallback.go)
|
||||
// replaced by ListResourcesClient
|
||||
type NodeClient interface {
|
||||
ListNodes(ctx context.Context, req proto.ListNodesRequest) (nodes []types.Server, nextKey string, err error)
|
||||
}
|
||||
|
||||
// GetNodesWithLabels is a helper for getting a list of nodes with optional label-based filtering. In addition to
|
||||
// iterating pages, it also correctly handles downsizing pages when LimitExceeded errors are encountered.
|
||||
//
|
||||
// DELETE IN 11.0.0 replaced by GetResourcesWithFilters.
|
||||
func GetNodesWithLabels(ctx context.Context, clt NodeClient, namespace string, labels map[string]string) ([]types.Server, error) {
|
||||
// Retrieve the complete list of nodes in chunks.
|
||||
var (
|
||||
nodes []types.Server
|
||||
startKey string
|
||||
chunkSize = int32(defaults.DefaultChunkSize)
|
||||
)
|
||||
for {
|
||||
resp, nextKey, err := clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: namespace,
|
||||
Limit: chunkSize,
|
||||
StartKey: startKey,
|
||||
Labels: labels,
|
||||
})
|
||||
if trace.IsLimitExceeded(err) {
|
||||
// Cut chunkSize in half if gRPC max message size is exceeded.
|
||||
chunkSize = chunkSize / 2
|
||||
// This is an extremely unlikely scenario, but better to cover it anyways.
|
||||
if chunkSize == 0 {
|
||||
return nil, trace.Wrap(trail.FromGRPC(err), "Node is too large to retrieve over gRPC (over 4MiB).")
|
||||
}
|
||||
continue
|
||||
} else if err != nil {
|
||||
return nil, trail.FromGRPC(err)
|
||||
}
|
||||
|
||||
// perform client-side filtering in case we're dealing with an older auth server which
|
||||
// does not support server-side filtering.
|
||||
for _, node := range resp {
|
||||
if node.MatchAgainst(labels) {
|
||||
nodes = append(nodes, node)
|
||||
}
|
||||
}
|
||||
|
||||
startKey = nextKey
|
||||
if startKey == "" {
|
||||
return nodes, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ListNodes returns a paginated list of nodes that the user has access to in the given namespace.
|
||||
// nextKey can be used as startKey in another call to ListNodes to retrieve the next page of nodes.
|
||||
// ListNodes will return a trace.LimitExceeded error if the page of nodes retrieved exceeds 4MiB.
|
||||
func (c *Client) ListNodes(ctx context.Context, req proto.ListNodesRequest) ([]types.Server, string, error) {
|
||||
resp, err := c.ListResources(ctx, proto.ListResourcesRequest{
|
||||
ResourceType: types.KindNode,
|
||||
Namespace: req.Namespace,
|
||||
StartKey: req.StartKey,
|
||||
Limit: req.Limit,
|
||||
Labels: req.Labels,
|
||||
})
|
||||
if err != nil {
|
||||
if trace.IsNotImplemented(err) {
|
||||
return c.listNodesFallback(ctx, req)
|
||||
}
|
||||
|
||||
return nil, "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
servers := make([]types.Server, len(resp.Resources))
|
||||
for i, resource := range resp.Resources {
|
||||
server, ok := resource.(*types.ServerV2)
|
||||
if !ok {
|
||||
return nil, "", trace.BadParameter("invalid node type %T", resource)
|
||||
}
|
||||
|
||||
servers[i] = server
|
||||
}
|
||||
|
||||
return servers, resp.NextKey, nil
|
||||
}
|
||||
|
||||
// listNodesFallback previous implementation of `ListNodes` function using
|
||||
// `ListNodes` RPC call.
|
||||
// DELETE IN 10.0
|
||||
func (c *Client) listNodesFallback(ctx context.Context, req proto.ListNodesRequest) (nodes []types.Server, nextKey string, err error) {
|
||||
if req.Namespace == "" {
|
||||
return nil, "", trace.BadParameter("missing parameter namespace")
|
||||
}
|
||||
if req.Limit <= 0 {
|
||||
return nil, "", trace.BadParameter("nonpositive parameter limit")
|
||||
}
|
||||
|
||||
resp, err := c.grpc.ListNodes(ctx, &req, c.callOpts...)
|
||||
if err != nil {
|
||||
return nil, "", trail.FromGRPC(err)
|
||||
}
|
||||
|
||||
nodes = make([]types.Server, len(resp.Servers))
|
||||
for i, node := range resp.Servers {
|
||||
nodes[i] = node
|
||||
}
|
||||
|
||||
return nodes, resp.NextKey, nil
|
||||
}
|
||||
|
||||
// UpsertNode is used by SSH servers to report their presence
|
||||
// to the auth servers in form of heartbeat expiring after ttl period.
|
||||
func (c *Client) UpsertNode(ctx context.Context, node types.Server) (*types.KeepAlive, error) {
|
||||
|
||||
@@ -482,43 +482,6 @@ func TestWaitForConnectionReady(t *testing.T) {
|
||||
require.Error(t, clt.waitForConnectionReady(ctx))
|
||||
}
|
||||
|
||||
func TestLimitExceeded(t *testing.T) {
|
||||
t.Parallel()
|
||||
ctx := context.Background()
|
||||
srv := startMockServer(t)
|
||||
|
||||
// Create client
|
||||
clt, err := srv.NewClient(ctx)
|
||||
require.NoError(t, err)
|
||||
|
||||
// ListNodes should return a limit exceeded error when exceeding gRPC message size limit.
|
||||
_, _, err = clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: defaults.Namespace,
|
||||
Limit: 50,
|
||||
})
|
||||
require.IsType(t, &trace.LimitExceededError{}, err.(*trace.TraceErr).OrigError())
|
||||
|
||||
// GetNodes should retrieve all nodes and transparently handle limit exceeded errors.
|
||||
expectedResources, err := testResources(types.KindNode, defaults.Namespace)
|
||||
require.NoError(t, err)
|
||||
|
||||
expectedNodes := make([]types.Server, len(expectedResources))
|
||||
for i, expectedResource := range expectedResources {
|
||||
var ok bool
|
||||
expectedNodes[i], ok = expectedResource.(*types.ServerV2)
|
||||
require.True(t, ok)
|
||||
}
|
||||
|
||||
resp, err := clt.GetNodes(ctx, defaults.Namespace)
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, expectedNodes, resp)
|
||||
|
||||
// GetNodes should fail with a limit exceeded error if a
|
||||
// single node is too big to send over gRPC (over 4MB).
|
||||
_, err = clt.GetNodes(ctx, fiveMBNode)
|
||||
require.IsType(t, &trace.LimitExceededError{}, err.(*trace.TraceErr).OrigError())
|
||||
}
|
||||
|
||||
func TestListResources(t *testing.T) {
|
||||
t.Parallel()
|
||||
ctx := context.Background()
|
||||
|
||||
+653
-1423
File diff suppressed because it is too large
Load Diff
@@ -1159,28 +1159,6 @@ message Events {
|
||||
string LastKey = 2;
|
||||
}
|
||||
|
||||
// DELETE IN 10.0
|
||||
message ListNodesRequest {
|
||||
// Namespace is the namespace of resources.
|
||||
string Namespace = 1;
|
||||
// Limit is the maximum amount of nodes to retrieve.
|
||||
int32 Limit = 2;
|
||||
// StartKey is used to start listing nodes from a specific spot. This should
|
||||
// be set to the previous NextKey value if using pagination, or left empty.
|
||||
string StartKey = 3;
|
||||
// Labels is a label-based matcher if non-empty.
|
||||
map<string, string> Labels = 4;
|
||||
}
|
||||
|
||||
// DELETE IN 10.0
|
||||
message ListNodesResponse {
|
||||
// Servers is a list of servers.
|
||||
repeated types.ServerV2 Servers = 1;
|
||||
// NextKey is the next Key to use as StartKey in a ListNodesRequest to continue
|
||||
// retrieving pages of nodes. If NextKey is empty, there are no more pages.
|
||||
string NextKey = 2;
|
||||
}
|
||||
|
||||
message GetLocksRequest {
|
||||
// Targets is a list of targets. Every returned lock must match at least
|
||||
// one of the targets.
|
||||
@@ -1500,6 +1478,7 @@ message PaginatedResource {
|
||||
// one type of resource can be retrieved per request.
|
||||
message ListResourcesRequest {
|
||||
// ResourceType is the resource that is going to be retrieved.
|
||||
// This only needs to be set explicitly for the `ListResources` rpc.
|
||||
string ResourceType = 1 [ (gogoproto.jsontag) = "resource_type,omitempty" ];
|
||||
// Namespace is the namespace of resources.
|
||||
string Namespace = 2 [ (gogoproto.jsontag) = "namespace,omitempty" ];
|
||||
@@ -1720,12 +1699,6 @@ service AuthService {
|
||||
|
||||
// GetNode retrieves a node described by the given request.
|
||||
rpc GetNode(types.ResourceInNamespaceRequest) returns (types.ServerV2);
|
||||
// GetNodes retrieves all nodes.
|
||||
// DELETE IN 8.0.0 in favor of ListNodes
|
||||
rpc GetNodes(types.ResourcesInNamespaceRequest) returns (types.ServerV2List);
|
||||
// ListNodes retrieves a paginated list of nodes.
|
||||
// DELETE IN 10.0. Deprecated, use ListResources.
|
||||
rpc ListNodes(ListNodesRequest) returns (ListNodesResponse) { option deprecated = true; };
|
||||
// UpsertNode upserts a node in a backend.
|
||||
rpc UpsertNode(types.ServerV2) returns (types.KeepAlive);
|
||||
// DeleteNode deletes an existing node in a backend described by the given request.
|
||||
|
||||
+1
-1
Submodule e updated: d2c63fbb62...d2f3597890
+3
-6
@@ -201,7 +201,7 @@ type ReadProxyAccessPoint interface {
|
||||
GetNode(ctx context.Context, namespace, name string) (types.Server, error)
|
||||
|
||||
// GetNodes returns a list of registered servers for this cluster.
|
||||
GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error)
|
||||
GetNodes(ctx context.Context, namespace string) ([]types.Server, error)
|
||||
|
||||
// GetProxies returns a list of proxy servers registered in the cluster
|
||||
GetProxies() ([]types.Server, error)
|
||||
@@ -329,7 +329,7 @@ type ReadRemoteProxyAccessPoint interface {
|
||||
GetNode(ctx context.Context, namespace, name string) (types.Server, error)
|
||||
|
||||
// GetNodes returns a list of registered servers for this cluster.
|
||||
GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error)
|
||||
GetNodes(ctx context.Context, namespace string) ([]types.Server, error)
|
||||
|
||||
// GetProxies returns a list of proxy servers registered in the cluster
|
||||
GetProxies() ([]types.Server, error)
|
||||
@@ -697,10 +697,7 @@ type Cache interface {
|
||||
GetNode(ctx context.Context, namespace, name string) (types.Server, error)
|
||||
|
||||
// GetNodes returns a list of registered servers for this cluster.
|
||||
GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error)
|
||||
|
||||
// ListNodes returns a paginated list of registered servers for this cluster.
|
||||
ListNodes(ctx context.Context, req proto.ListNodesRequest) (nodes []types.Server, nextKey string, err error)
|
||||
GetNodes(ctx context.Context, namespace string) ([]types.Server, error)
|
||||
|
||||
// GetProxies returns a list of proxy servers registered in the cluster
|
||||
GetProxies() ([]types.Server, error)
|
||||
|
||||
@@ -121,12 +121,7 @@ func NewAPIServer(config *APIConfig) (http.Handler, error) {
|
||||
srv.DELETE("/:version/users/:user/web/sessions/:sid", srv.withAuth(srv.deleteWebSession))
|
||||
|
||||
// Servers and presence heartbeat
|
||||
srv.POST("/:version/namespaces/:namespace/nodes", srv.withAuth(srv.upsertNode))
|
||||
srv.POST("/:version/namespaces/:namespace/nodes/keepalive", srv.withAuth(srv.keepAliveNode))
|
||||
srv.PUT("/:version/namespaces/:namespace/nodes", srv.withAuth(srv.upsertNodes))
|
||||
srv.GET("/:version/namespaces/:namespace/nodes", srv.withAuth(srv.getNodes))
|
||||
srv.DELETE("/:version/namespaces/:namespace/nodes", srv.withAuth(srv.deleteAllNodes))
|
||||
srv.DELETE("/:version/namespaces/:namespace/nodes/:name", srv.withAuth(srv.deleteNode))
|
||||
srv.POST("/:version/authservers", srv.withAuth(srv.upsertAuthServer))
|
||||
srv.GET("/:version/authservers", srv.withAuth(srv.getAuthServers))
|
||||
srv.POST("/:version/proxies", srv.withAuth(srv.upsertProxy))
|
||||
@@ -350,84 +345,6 @@ func (s *APIServer) keepAliveNode(auth ClientI, w http.ResponseWriter, r *http.R
|
||||
return message("ok"), nil
|
||||
}
|
||||
|
||||
type upsertNodesReq struct {
|
||||
Nodes json.RawMessage `json:"nodes"`
|
||||
Namespace string `json:"namespace"`
|
||||
}
|
||||
|
||||
// upsertNodes is used to bulk insert nodes into the backend.
|
||||
func (s *APIServer) upsertNodes(auth ClientI, w http.ResponseWriter, r *http.Request, p httprouter.Params, version string) (interface{}, error) {
|
||||
var req upsertNodesReq
|
||||
if err := httplib.ReadJSON(r, &req); err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
if !types.IsValidNamespace(req.Namespace) {
|
||||
return nil, trace.BadParameter("invalid namespace %q", req.Namespace)
|
||||
}
|
||||
|
||||
nodes, err := services.UnmarshalServers(req.Nodes)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
err = auth.UpsertNodes(req.Namespace, nodes)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
return message("ok"), nil
|
||||
}
|
||||
|
||||
// upsertNode is called by remote SSH nodes when they ping back into the auth service
|
||||
func (s *APIServer) upsertNode(auth ClientI, w http.ResponseWriter, r *http.Request, p httprouter.Params, version string) (interface{}, error) {
|
||||
return s.upsertServer(auth, types.RoleNode, r, p)
|
||||
}
|
||||
|
||||
// getNodes returns registered SSH nodes
|
||||
func (s *APIServer) getNodes(auth ClientI, w http.ResponseWriter, r *http.Request, p httprouter.Params, version string) (interface{}, error) {
|
||||
namespace := p.ByName("namespace")
|
||||
if !types.IsValidNamespace(namespace) {
|
||||
return nil, trace.BadParameter("invalid namespace %q", namespace)
|
||||
}
|
||||
|
||||
servers, err := auth.GetNodes(r.Context(), namespace)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
return marshalServers(servers, version)
|
||||
}
|
||||
|
||||
// deleteAllNodes deletes all nodes
|
||||
func (s *APIServer) deleteAllNodes(auth ClientI, w http.ResponseWriter, r *http.Request, p httprouter.Params, version string) (interface{}, error) {
|
||||
namespace := p.ByName("namespace")
|
||||
if !types.IsValidNamespace(namespace) {
|
||||
return nil, trace.BadParameter("invalid namespace %q", namespace)
|
||||
}
|
||||
err := auth.DeleteAllNodes(r.Context(), namespace)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
return message("ok"), nil
|
||||
}
|
||||
|
||||
// deleteNode deletes node
|
||||
func (s *APIServer) deleteNode(auth ClientI, w http.ResponseWriter, r *http.Request, p httprouter.Params, version string) (interface{}, error) {
|
||||
namespace := p.ByName("namespace")
|
||||
if !types.IsValidNamespace(namespace) {
|
||||
return nil, trace.BadParameter("invalid namespace %q", namespace)
|
||||
}
|
||||
name := p.ByName("name")
|
||||
if name == "" {
|
||||
return nil, trace.BadParameter("missing node name")
|
||||
}
|
||||
err := auth.DeleteNode(r.Context(), namespace, name)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
return message("ok"), nil
|
||||
}
|
||||
|
||||
// upsertProxy is called by remote SSH nodes when they ping back into the auth service
|
||||
func (s *APIServer) upsertProxy(auth ClientI, w http.ResponseWriter, r *http.Request, p httprouter.Params, version string) (interface{}, error) {
|
||||
return s.upsertServer(auth, types.RoleProxy, r, p)
|
||||
|
||||
+2
-33
@@ -2739,39 +2739,8 @@ func (a *Server) GetNamespaces() ([]types.Namespace, error) {
|
||||
}
|
||||
|
||||
// GetNodes returns nodes from the cache
|
||||
func (a *Server) GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
return a.GetCache().GetNodes(ctx, namespace, opts...)
|
||||
}
|
||||
|
||||
// ListNodes lists nodes from the cache
|
||||
func (a *Server) ListNodes(ctx context.Context, req proto.ListNodesRequest) ([]types.Server, string, error) {
|
||||
return a.GetCache().ListNodes(ctx, req)
|
||||
}
|
||||
|
||||
// NodePageFunc is a function to run on each page iterated over.
|
||||
type NodePageFunc func(next []types.Server) (stop bool, err error)
|
||||
|
||||
// IterateNodePages can be used to iterate over pages of nodes.
|
||||
func (a *Server) IterateNodePages(ctx context.Context, req proto.ListNodesRequest, f NodePageFunc) (string, error) {
|
||||
for {
|
||||
nextPage, nextKey, err := a.ListNodes(ctx, req)
|
||||
if err != nil {
|
||||
return "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
stop, err := f(nextPage)
|
||||
if err != nil {
|
||||
return "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
// Iterator stopped before end of pages or
|
||||
// there are no more pages, return nextKey
|
||||
if stop || nextKey == "" {
|
||||
return nextKey, nil
|
||||
}
|
||||
|
||||
req.StartKey = nextKey
|
||||
}
|
||||
func (a *Server) GetNodes(ctx context.Context, namespace string) ([]types.Server, error) {
|
||||
return a.GetCache().GetNodes(ctx, namespace)
|
||||
}
|
||||
|
||||
// ResourcePageFunc is a function to run on each page iterated over.
|
||||
|
||||
@@ -32,7 +32,6 @@ import (
|
||||
apievents "github.com/gravitational/teleport/api/types/events"
|
||||
"github.com/gravitational/teleport/api/types/wrappers"
|
||||
apiutils "github.com/gravitational/teleport/api/utils"
|
||||
"github.com/gravitational/teleport/lib/backend"
|
||||
"github.com/gravitational/teleport/lib/defaults"
|
||||
"github.com/gravitational/teleport/lib/events"
|
||||
"github.com/gravitational/teleport/lib/modules"
|
||||
@@ -623,14 +622,6 @@ func (a *ServerWithRoles) GenerateHostCerts(ctx context.Context, req *proto.Host
|
||||
return a.authServer.GenerateHostCerts(ctx, req)
|
||||
}
|
||||
|
||||
// UpsertNodes bulk upserts nodes into the backend.
|
||||
func (a *ServerWithRoles) UpsertNodes(namespace string, servers []types.Server) error {
|
||||
if err := a.action(namespace, types.KindNode, types.VerbCreate, types.VerbUpdate); err != nil {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
return a.authServer.UpsertNodes(namespace, servers)
|
||||
}
|
||||
|
||||
func (a *ServerWithRoles) UpsertNode(ctx context.Context, s types.Server) (*types.KeepAlive, error) {
|
||||
if err := a.action(s.GetNamespace(), types.KindNode, types.VerbCreate, types.VerbUpdate); err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
@@ -875,14 +866,14 @@ func (a *ServerWithRoles) GetNode(ctx context.Context, namespace, name string) (
|
||||
return node, nil
|
||||
}
|
||||
|
||||
func (a *ServerWithRoles) GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
func (a *ServerWithRoles) GetNodes(ctx context.Context, namespace string) ([]types.Server, error) {
|
||||
if err := a.action(namespace, types.KindNode, types.VerbList); err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
// Fetch full list of nodes in the backend.
|
||||
startFetch := time.Now()
|
||||
nodes, err := a.authServer.GetNodes(ctx, namespace, opts...)
|
||||
nodes, err := a.authServer.GetNodes(ctx, namespace)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
@@ -1296,66 +1287,6 @@ func (a *ServerWithRoles) ListWindowsDesktops(ctx context.Context, req types.Lis
|
||||
return nil, trace.NotImplemented(notImplementedMessage)
|
||||
}
|
||||
|
||||
// ListNodes returns a paginated list of nodes filtered by user access.
|
||||
//
|
||||
// DELETE IN 10.0.0 in favor of ListResources.
|
||||
func (a *ServerWithRoles) ListNodes(ctx context.Context, req proto.ListNodesRequest) (page []types.Server, nextKey string, err error) {
|
||||
if err := a.action(req.Namespace, types.KindNode, types.VerbList); err != nil {
|
||||
return nil, "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
return a.filterAndListNodes(ctx, req)
|
||||
}
|
||||
|
||||
// DELETE IN 10.0.0 in favor of ListResources.
|
||||
func (a *ServerWithRoles) filterAndListNodes(ctx context.Context, req proto.ListNodesRequest) (page []types.Server, nextKey string, err error) {
|
||||
limit := int(req.Limit)
|
||||
if limit <= 0 {
|
||||
return nil, "", trace.BadParameter("nonpositive parameter limit")
|
||||
}
|
||||
|
||||
// move labels out of request so that we can perform label-based filtering *after* RBAC filtering.
|
||||
realLabels := req.Labels
|
||||
req.Labels = nil
|
||||
|
||||
checker, err := newNodeChecker(a.context, a.authServer)
|
||||
if err != nil {
|
||||
return nil, "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
page = make([]types.Server, 0, limit)
|
||||
nextKey, err = a.authServer.IterateNodePages(ctx, req, func(nextPage []types.Server) (bool, error) {
|
||||
// Retrieve and filter pages of nodes until we can fill a page or run out of nodes.
|
||||
filteredPage, err := a.filterNodes(checker, nextPage)
|
||||
if err != nil {
|
||||
return false, trace.Wrap(err)
|
||||
}
|
||||
|
||||
// add all matching nodes to page
|
||||
for _, node := range filteredPage {
|
||||
if len(page) == limit {
|
||||
// page is filled, stop processing
|
||||
break
|
||||
}
|
||||
if node.MatchAgainst(realLabels) {
|
||||
page = append(page, node)
|
||||
}
|
||||
}
|
||||
|
||||
return len(page) == limit, nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
// Filled a page, reset nextKey in case the last node was cut out.
|
||||
if len(page) == limit {
|
||||
nextKey = backend.NextPaginationKey(page[len(page)-1])
|
||||
}
|
||||
|
||||
return page, nextKey, nil
|
||||
}
|
||||
|
||||
func (a *ServerWithRoles) UpsertAuthServer(s types.Server) error {
|
||||
if err := a.action(apidefaults.Namespace, types.KindAuthServer, types.VerbCreate, types.VerbUpdate); err != nil {
|
||||
return trace.Wrap(err)
|
||||
|
||||
@@ -1395,67 +1395,6 @@ func TestSessionRecordingConfigRBAC(t *testing.T) {
|
||||
})
|
||||
}
|
||||
|
||||
// TestListNodes users can retrieve nodes with the appropriate permissions.
|
||||
func TestListNodes(t *testing.T) {
|
||||
t.Parallel()
|
||||
ctx := context.Background()
|
||||
srv := newTestTLSServer(t)
|
||||
|
||||
// Create test nodes.
|
||||
for i := 0; i < 10; i++ {
|
||||
name := uuid.New().String()
|
||||
node, err := types.NewServerWithLabels(
|
||||
name,
|
||||
types.KindNode,
|
||||
types.ServerSpecV2{},
|
||||
map[string]string{"name": name},
|
||||
)
|
||||
require.NoError(t, err)
|
||||
|
||||
_, err = srv.Auth().UpsertNode(ctx, node)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
testNodes, err := srv.Auth().GetNodes(ctx, defaults.Namespace)
|
||||
require.NoError(t, err)
|
||||
|
||||
// create user, role, and client
|
||||
username := "user"
|
||||
user, role, err := CreateUserAndRole(srv.Auth(), username, nil)
|
||||
require.NoError(t, err)
|
||||
identity := TestUser(user.GetName())
|
||||
clt, err := srv.NewClient(identity)
|
||||
require.NoError(t, err)
|
||||
|
||||
// permit user to list all nodes
|
||||
role.SetNodeLabels(types.Allow, types.Labels{types.Wildcard: {types.Wildcard}})
|
||||
require.NoError(t, srv.Auth().UpsertRole(ctx, role))
|
||||
|
||||
// listing nodes 0-4 should list first 5 nodes
|
||||
nodes, _, err := clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: defaults.Namespace,
|
||||
Limit: 5,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 5, len(nodes))
|
||||
expectedNodes := testNodes[:5]
|
||||
require.Empty(t, cmp.Diff(expectedNodes, nodes))
|
||||
|
||||
// remove permission for third node
|
||||
role.SetNodeLabels(types.Deny, types.Labels{"name": {testNodes[3].GetName()}})
|
||||
require.NoError(t, srv.Auth().UpsertRole(ctx, role))
|
||||
|
||||
// listing nodes 0-4 should skip the third node and add the fifth to the end.
|
||||
nodes, _, err = clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: defaults.Namespace,
|
||||
Limit: 5,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 5, len(nodes))
|
||||
expectedNodes = append(testNodes[:3], testNodes[4:6]...)
|
||||
require.Empty(t, cmp.Diff(expectedNodes, nodes))
|
||||
}
|
||||
|
||||
// TestGetAndList_Nodes users can retrieve nodes with various filters
|
||||
// and with the appropriate permissions.
|
||||
func TestGetAndList_Nodes(t *testing.T) {
|
||||
|
||||
@@ -530,24 +530,6 @@ func (c *Client) KeepAliveServer(ctx context.Context, keepAlive types.KeepAlive)
|
||||
return trace.BadParameter("not implemented, use StreamKeepAlives instead")
|
||||
}
|
||||
|
||||
// UpsertNodes bulk inserts nodes.
|
||||
func (c *Client) UpsertNodes(namespace string, servers []types.Server) error {
|
||||
if namespace == "" {
|
||||
return trace.BadParameter("missing node namespace")
|
||||
}
|
||||
|
||||
bytes, err := services.MarshalServers(servers)
|
||||
if err != nil {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
args := &upsertNodesReq{
|
||||
Namespace: namespace,
|
||||
Nodes: bytes,
|
||||
}
|
||||
_, err = c.PutJSON(context.TODO(), c.Endpoint("namespaces", namespace, "nodes"), args)
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
|
||||
// UpsertReverseTunnel is used by admins to create a new reverse tunnel
|
||||
// to the remote proxy to bypass firewall restrictions
|
||||
func (c *Client) UpsertReverseTunnel(tunnel types.ReverseTunnel) error {
|
||||
|
||||
@@ -2586,52 +2586,6 @@ func (g *GRPCServer) GetNode(ctx context.Context, req *types.ResourceInNamespace
|
||||
return serverV2, nil
|
||||
}
|
||||
|
||||
// GetNodes retrieves all nodes in the given namespace.
|
||||
// DELETE IN 8.0.0 in favor of ListNodes
|
||||
func (g *GRPCServer) GetNodes(ctx context.Context, req *types.ResourcesInNamespaceRequest) (*types.ServerV2List, error) {
|
||||
auth, err := g.authenticate(ctx)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
ns, err := auth.ServerWithRoles.GetNodes(ctx, req.Namespace)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
serversV2 := make([]*types.ServerV2, len(ns))
|
||||
for i, t := range ns {
|
||||
var ok bool
|
||||
if serversV2[i], ok = t.(*types.ServerV2); !ok {
|
||||
return nil, trace.Errorf("encountered unexpected node type: %T", t)
|
||||
}
|
||||
}
|
||||
return &types.ServerV2List{
|
||||
Servers: serversV2,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// ListNodes retrieves a paginated list of nodes in the given namespace.
|
||||
func (g *GRPCServer) ListNodes(ctx context.Context, req *proto.ListNodesRequest) (*proto.ListNodesResponse, error) {
|
||||
auth, err := g.authenticate(ctx)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
ns, nextKey, err := auth.ServerWithRoles.ListNodes(ctx, *req)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
serversV2 := make([]*types.ServerV2, len(ns))
|
||||
for i, t := range ns {
|
||||
var ok bool
|
||||
if serversV2[i], ok = t.(*types.ServerV2); !ok {
|
||||
return nil, trace.Errorf("encountered unexpected node type: %T", t)
|
||||
}
|
||||
}
|
||||
return &proto.ListNodesResponse{
|
||||
Servers: serversV2,
|
||||
NextKey: nextKey,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// UpsertNode upserts a node.
|
||||
func (g *GRPCServer) UpsertNode(ctx context.Context, node *types.ServerV2) (*types.KeepAlive, error) {
|
||||
auth, err := g.authenticate(ctx)
|
||||
|
||||
@@ -39,7 +39,6 @@ import (
|
||||
"github.com/gravitational/teleport/lib/auth/mocku2f"
|
||||
"github.com/gravitational/teleport/lib/auth/native"
|
||||
wanlib "github.com/gravitational/teleport/lib/auth/webauthn"
|
||||
"github.com/gravitational/teleport/lib/backend"
|
||||
"github.com/gravitational/teleport/lib/defaults"
|
||||
"github.com/gravitational/teleport/lib/services"
|
||||
"github.com/gravitational/teleport/lib/tlsca"
|
||||
@@ -1391,51 +1390,6 @@ func TestNodesCRUD(t *testing.T) {
|
||||
|
||||
// Run NodeGetters in nested subtests to allow parallelization.
|
||||
t.Run("NodeGetters", func(t *testing.T) {
|
||||
t.Run("List Nodes", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
// List nodes one at a time, last page should be empty.
|
||||
|
||||
// First node.
|
||||
nodes, nextKey, err := clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: 1,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, nodes, 1)
|
||||
require.Empty(t, cmp.Diff([]types.Server{node1}, nodes,
|
||||
cmpopts.IgnoreFields(types.Metadata{}, "ID")))
|
||||
require.Equal(t, backend.NextPaginationKey(node1), nextKey)
|
||||
|
||||
// Second node (last).
|
||||
nodes, nextKey, err = clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: 1,
|
||||
StartKey: nextKey,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Len(t, nodes, 1)
|
||||
require.Empty(t, cmp.Diff([]types.Server{node2}, nodes,
|
||||
cmpopts.IgnoreFields(types.Metadata{}, "ID")))
|
||||
require.Empty(t, nextKey)
|
||||
|
||||
// ListNodes should not fail if namespace is empty
|
||||
_, _, err = clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Limit: 1,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
// ListNodes should fail if limit is nonpositive
|
||||
_, _, err = clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
})
|
||||
require.IsType(t, &trace.BadParameterError{}, err.(*trace.TraceErr).OrigError())
|
||||
|
||||
_, _, err = clt.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: -1,
|
||||
})
|
||||
require.IsType(t, &trace.BadParameterError{}, err.(*trace.TraceErr).OrigError())
|
||||
})
|
||||
t.Run("GetNodes", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
// Get all nodes
|
||||
|
||||
@@ -21,121 +21,13 @@ import (
|
||||
"encoding/json"
|
||||
"net/url"
|
||||
|
||||
"github.com/gravitational/teleport/api/client"
|
||||
"github.com/gravitational/teleport/api/client/proto"
|
||||
"github.com/gravitational/teleport/api/types"
|
||||
"github.com/gravitational/teleport/lib/services"
|
||||
|
||||
"github.com/gravitational/trace"
|
||||
)
|
||||
|
||||
// httpfallback.go holds endpoints that have been converted to gRPC
|
||||
// but still need http fallback logic in the old client.
|
||||
|
||||
// DeleteAllNodes deletes all nodes in a given namespace
|
||||
func (c *Client) DeleteAllNodes(ctx context.Context, namespace string) error {
|
||||
if err := c.APIClient.DeleteAllNodes(ctx, namespace); err != nil {
|
||||
if !trace.IsNotImplemented(err) {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
} else {
|
||||
return nil
|
||||
}
|
||||
|
||||
_, err := c.Delete(ctx, c.Endpoint("namespaces", namespace, "nodes"))
|
||||
if err != nil {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeleteNode deletes node in the namespace by name
|
||||
func (c *Client) DeleteNode(ctx context.Context, namespace string, name string) error {
|
||||
if err := c.APIClient.DeleteNode(ctx, namespace, name); err != nil {
|
||||
if !trace.IsNotImplemented(err) {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
} else {
|
||||
return nil
|
||||
}
|
||||
|
||||
_, err := c.Delete(ctx, c.Endpoint("namespaces", namespace, "nodes", name))
|
||||
if err != nil {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type nodeClient interface {
|
||||
ListNodes(ctx context.Context, req proto.ListNodesRequest) (nodes []types.Server, nextKey string, err error)
|
||||
GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error)
|
||||
}
|
||||
|
||||
// GetNodesWithLabels is a helper for getting a list of nodes with optional label-based filtering. This is essentially
|
||||
// a wrapper around client.GetNodesWithLabels that performs fallback on NotImplemented errors.
|
||||
//
|
||||
// DELETE IN 11.0.0, this function is only called by lib/client/client.go (*ProxyClient).FindServersByLabels
|
||||
// which is also marked for deletion (replaced by FindNodesByFilters).
|
||||
func GetNodesWithLabels(ctx context.Context, clt nodeClient, namespace string, labels map[string]string) ([]types.Server, error) {
|
||||
nodes, err := client.GetNodesWithLabels(ctx, clt, namespace, labels)
|
||||
if err == nil || !trace.IsNotImplemented(err) {
|
||||
return nodes, trace.Wrap(err)
|
||||
}
|
||||
|
||||
nodes, err = clt.GetNodes(ctx, namespace)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
var filtered []types.Server
|
||||
|
||||
// we had to fallback to a method that does not perform server-side filtering,
|
||||
// so filter here instead.
|
||||
for _, node := range nodes {
|
||||
if node.MatchAgainst(labels) {
|
||||
filtered = append(filtered, node)
|
||||
}
|
||||
}
|
||||
|
||||
return filtered, nil
|
||||
}
|
||||
|
||||
// GetNodes returns the list of servers registered in the cluster.
|
||||
//
|
||||
// DELETE IN 11.0.0, replaced by GetResourcesWithFilters
|
||||
func (c *Client) GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
if resp, err := c.APIClient.GetNodes(ctx, namespace); err != nil {
|
||||
if !trace.IsNotImplemented(err) {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
} else {
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
out, err := c.Get(ctx, c.Endpoint("namespaces", namespace, "nodes"), url.Values{})
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
var items []json.RawMessage
|
||||
if err := json.Unmarshal(out.Bytes(), &items); err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
re := make([]types.Server, len(items))
|
||||
for i, raw := range items {
|
||||
s, err := services.UnmarshalServer(
|
||||
raw,
|
||||
types.KindNode,
|
||||
opts...)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
re[i] = s
|
||||
}
|
||||
|
||||
return re, nil
|
||||
}
|
||||
|
||||
// GetDomainName returns local auth domain of the current auth server
|
||||
// DELETE IN 11.0.0
|
||||
func (c *Client) GetDomainName(ctx context.Context) (string, error) {
|
||||
|
||||
Vendored
+4
-72
@@ -1412,7 +1412,7 @@ type getNodesCacheKey struct {
|
||||
var _ map[getNodesCacheKey]struct{} // compile-time hashability check
|
||||
|
||||
// GetNodes is a part of auth.Cache implementation
|
||||
func (c *Cache) GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
func (c *Cache) GetNodes(ctx context.Context, namespace string) ([]types.Server, error) {
|
||||
rg, err := c.read()
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
@@ -1420,7 +1420,7 @@ func (c *Cache) GetNodes(ctx context.Context, namespace string, opts ...services
|
||||
defer rg.Release()
|
||||
|
||||
if !rg.IsCacheRead() {
|
||||
cachedNodes, err := c.getNodesWithTTLCache(ctx, rg, namespace, opts...)
|
||||
cachedNodes, err := c.getNodesWithTTLCache(ctx, rg, namespace)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
@@ -1431,7 +1431,7 @@ func (c *Cache) GetNodes(ctx context.Context, namespace string, opts ...services
|
||||
return nodes, nil
|
||||
}
|
||||
|
||||
return rg.presence.GetNodes(ctx, namespace, opts...)
|
||||
return rg.presence.GetNodes(ctx, namespace)
|
||||
}
|
||||
|
||||
// getNodesWithTTLCache implements TTL-based caching for the GetNodes endpoint. All nodes that will be returned from the caching layer
|
||||
@@ -1441,7 +1441,7 @@ func (c *Cache) getNodesWithTTLCache(ctx context.Context, rg readGuard, namespac
|
||||
ni, err := c.fnCache.Get(ctx, getNodesCacheKey{namespace}, func(ctx context.Context) (interface{}, error) {
|
||||
// use cache's close context instead of request context in order to ensure
|
||||
// that we don't cache a context cancellation error.
|
||||
nodes, err := rg.presence.GetNodes(c.ctx, namespace, opts...)
|
||||
nodes, err := rg.presence.GetNodes(c.ctx, namespace)
|
||||
ta(nodes)
|
||||
return nodes, err
|
||||
})
|
||||
@@ -1456,74 +1456,6 @@ func (c *Cache) getNodesWithTTLCache(ctx context.Context, rg readGuard, namespac
|
||||
return cachedNodes, nil
|
||||
}
|
||||
|
||||
// ListNodes is a part of auth.Cache implementation
|
||||
//
|
||||
// DELETE IN 10.0.0 in favor of ListResources.
|
||||
func (c *Cache) ListNodes(ctx context.Context, req proto.ListNodesRequest) ([]types.Server, string, error) {
|
||||
rg, err := c.read()
|
||||
if err != nil {
|
||||
return nil, "", trace.Wrap(err)
|
||||
}
|
||||
defer rg.Release()
|
||||
|
||||
if rg.IsCacheRead() {
|
||||
// cache is healthy, delegate to the standard ListNodes implementation
|
||||
return rg.presence.ListNodes(ctx, req)
|
||||
}
|
||||
|
||||
// Cache is not healthy. List nodes using TTL cache.
|
||||
return c.listNodesFromTTLCache(ctx, rg, req.Namespace, req.StartKey, req.Labels, int(req.Limit))
|
||||
}
|
||||
|
||||
// listNodesFromTTLCache Used when the cache is not healthy. It takes advantage
|
||||
// of TTL-based caching rather than caching individual page calls (very messy).
|
||||
// It relies on caching the result of the `GetNodes` endpoint and then "faking"
|
||||
// pagination.
|
||||
//
|
||||
// DELETE IN 10.0.0 in favor of ListResources.
|
||||
func (c *Cache) listNodesFromTTLCache(ctx context.Context, rg readGuard, namespace, startKey string, labels map[string]string, limit int) ([]types.Server, string, error) {
|
||||
if limit <= 0 {
|
||||
return nil, "", trace.BadParameter("nonpositive limit value")
|
||||
}
|
||||
|
||||
cachedNodes, err := c.getNodesWithTTLCache(ctx, rg, namespace)
|
||||
if err != nil {
|
||||
return nil, "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
// trim nodes that precede start key
|
||||
if startKey != "" {
|
||||
pageStart := 0
|
||||
for i, node := range cachedNodes {
|
||||
if node.GetName() < startKey {
|
||||
pageStart = i + 1
|
||||
} else {
|
||||
break
|
||||
}
|
||||
}
|
||||
cachedNodes = cachedNodes[pageStart:]
|
||||
}
|
||||
|
||||
// iterate and filter nodes, halting when we reach page limit
|
||||
var filtered []types.Server
|
||||
for _, node := range cachedNodes {
|
||||
if len(filtered) == limit {
|
||||
break
|
||||
}
|
||||
|
||||
if node.MatchAgainst(labels) {
|
||||
filtered = append(filtered, node.DeepCopy())
|
||||
}
|
||||
}
|
||||
|
||||
var nextKey string
|
||||
if len(filtered) == limit {
|
||||
nextKey = backend.NextPaginationKey(filtered[len(filtered)-1])
|
||||
}
|
||||
|
||||
return filtered, nextKey, nil
|
||||
}
|
||||
|
||||
// GetAuthServers returns a list of registered servers
|
||||
func (c *Cache) GetAuthServers() ([]types.Server, error) {
|
||||
rg, err := c.read()
|
||||
|
||||
Vendored
-82
@@ -737,62 +737,6 @@ func benchGetNodes(b *testing.B, nodeCount int) {
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
goos: linux
|
||||
goarch: amd64
|
||||
pkg: github.com/gravitational/teleport/lib/cache
|
||||
cpu: Intel(R) Core(TM) i9-10885H CPU @ 2.40GHz
|
||||
BenchmarkListMaxNodes-16 1 1136071399 ns/op
|
||||
*/
|
||||
func BenchmarkListMaxNodes(b *testing.B) {
|
||||
benchListNodes(b, backend.DefaultRangeLimit, apidefaults.DefaultChunkSize)
|
||||
}
|
||||
|
||||
func benchListNodes(b *testing.B, nodeCount int, pageSize int) {
|
||||
p, err := newPack(b.TempDir(), ForAuth, memoryBackend(true))
|
||||
require.NoError(b, err)
|
||||
defer p.Close()
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
for i := 0; i < nodeCount; i++ {
|
||||
func() {
|
||||
server := suite.NewServer(types.KindNode, uuid.New().String(), "127.0.0.1:2022", apidefaults.Namespace)
|
||||
_, err := p.presenceS.UpsertNode(ctx, server)
|
||||
require.NoError(b, err)
|
||||
timeout := time.NewTimer(time.Millisecond * 200)
|
||||
defer timeout.Stop()
|
||||
select {
|
||||
case event := <-p.eventsC:
|
||||
require.Equal(b, EventProcessed, event.Type)
|
||||
case <-timeout.C:
|
||||
b.Fatalf("timeout waiting for event, iteration=%d", i)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
b.ResetTimer()
|
||||
|
||||
for n := 0; n < b.N; n++ {
|
||||
var nodes []types.Server
|
||||
req := proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: int32(pageSize),
|
||||
}
|
||||
for {
|
||||
page, nextKey, err := p.cache.ListNodes(ctx, req)
|
||||
require.NoError(b, err)
|
||||
nodes = append(nodes, page...)
|
||||
require.True(b, len(page) == pageSize || nextKey == "")
|
||||
if nextKey == "" {
|
||||
break
|
||||
}
|
||||
req.StartKey = nextKey
|
||||
}
|
||||
require.Len(b, nodes, nodeCount)
|
||||
}
|
||||
}
|
||||
|
||||
// TestListResources_NodesTTLVariant verifies that the custom ListNodes impl that we fallback to when
|
||||
// using ttl-based caching works as expected.
|
||||
func TestListResources_NodesTTLVariant(t *testing.T) {
|
||||
@@ -844,32 +788,6 @@ func TestListResources_NodesTTLVariant(t *testing.T) {
|
||||
require.NoError(t, err)
|
||||
require.Len(t, allNodes, nodeCount)
|
||||
|
||||
// DELETE IN 10.0.0 this block with ListNodes is replaced
|
||||
// by the following block with ListResources test.
|
||||
var nodes []types.Server
|
||||
var startKey string
|
||||
for {
|
||||
page, nextKey, err := p.cache.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: int32(pageSize),
|
||||
StartKey: startKey,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
if nextKey != "" {
|
||||
require.Len(t, page, pageSize)
|
||||
}
|
||||
|
||||
nodes = append(nodes, page...)
|
||||
|
||||
startKey = nextKey
|
||||
|
||||
if startKey == "" {
|
||||
break
|
||||
}
|
||||
}
|
||||
require.Len(t, nodes, nodeCount)
|
||||
|
||||
var resources []types.ResourceWithLabels
|
||||
var listResourcesStartKey string
|
||||
sortBy := types.SortBy{
|
||||
|
||||
@@ -580,24 +580,6 @@ func (proxy *ProxyClient) isAuthBoring(ctx context.Context) (bool, error) {
|
||||
return resp.IsBoring, trace.Wrap(err)
|
||||
}
|
||||
|
||||
// FindServersByLabels returns list of the nodes which have labels exactly matching
|
||||
// the given label set.
|
||||
//
|
||||
// A server is matched when ALL labels match.
|
||||
// If no labels are passed, ALL nodes are returned.
|
||||
//
|
||||
// DELETE IN 11.0.0 replaced by FindNodesByFilters.
|
||||
func (proxy *ProxyClient) FindServersByLabels(ctx context.Context, namespace string, labels map[string]string) ([]types.Server, error) {
|
||||
if namespace == "" {
|
||||
return nil, trace.BadParameter(auth.MissingNamespaceError)
|
||||
}
|
||||
site, err := proxy.CurrentClusterAccessPoint(ctx, false)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
return auth.GetNodesWithLabels(ctx, site, namespace, labels)
|
||||
}
|
||||
|
||||
// FindServersByFilters returns list of the nodes which have filters matched.
|
||||
func (proxy *ProxyClient) FindNodesByFilters(ctx context.Context, req proto.ListResourcesRequest) ([]types.Server, error) {
|
||||
req.ResourceType = types.KindNode
|
||||
@@ -609,18 +591,6 @@ func (proxy *ProxyClient) FindNodesByFilters(ctx context.Context, req proto.List
|
||||
|
||||
resources, err := client.GetResourcesWithFilters(ctx, site, req)
|
||||
if err != nil {
|
||||
// ListResources for nodes not available, provide fallback.
|
||||
// Fallback does not support search/predicate support, so if users
|
||||
// provide them, it does nothing.
|
||||
//
|
||||
// DELETE IN 11.0.0
|
||||
if trace.IsNotImplemented(err) {
|
||||
servers, err := proxy.FindServersByLabels(ctx, req.Namespace, req.Labels)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
return servers, nil
|
||||
}
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
|
||||
@@ -64,7 +64,7 @@ func (c mockClient) GetRoles(ctx context.Context) ([]types.Role, error) {
|
||||
return c.roles, nil
|
||||
}
|
||||
|
||||
func (c mockClient) GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
func (c mockClient) GetNodes(ctx context.Context, namespace string) ([]types.Server, error) {
|
||||
return c.nodes, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -24,7 +24,6 @@ import (
|
||||
"github.com/gravitational/teleport/api/types"
|
||||
"github.com/gravitational/teleport/api/utils/sshutils"
|
||||
"github.com/gravitational/teleport/lib/auth"
|
||||
"github.com/gravitational/teleport/lib/services"
|
||||
"github.com/gravitational/trace"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
@@ -89,7 +88,7 @@ type mockLocalSiteClient struct {
|
||||
}
|
||||
|
||||
// called by (*localSite).sshTunnelStats() as part of (*localSite).periodicFunctions()
|
||||
func (mockLocalSiteClient) GetNodes(_ context.Context, _ string, _ ...services.MarshalOption) ([]types.Server, error) {
|
||||
func (mockLocalSiteClient) GetNodes(_ context.Context, _ string) ([]types.Server, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -118,11 +118,11 @@ func insertNodes(ctx context.Context, b *testing.B, svc services.Presence, nodeC
|
||||
}
|
||||
|
||||
// benchmarkGetNodes runs GetNodes b.N times.
|
||||
func benchmarkGetNodes(ctx context.Context, b *testing.B, svc services.Presence, nodeCount int, opts ...services.MarshalOption) {
|
||||
func benchmarkGetNodes(ctx context.Context, b *testing.B, svc services.Presence, nodeCount int) {
|
||||
var nodes []types.Server
|
||||
var err error
|
||||
for i := 0; i < b.N; i++ {
|
||||
nodes, err = svc.GetNodes(ctx, apidefaults.Namespace, opts...)
|
||||
nodes, err = svc.GetNodes(ctx, apidefaults.Namespace)
|
||||
require.NoError(b, err)
|
||||
}
|
||||
// do *something* with the loop result. probably unnecessary since the loop
|
||||
|
||||
@@ -206,7 +206,7 @@ func (s *PresenceService) GetNode(ctx context.Context, namespace, name string) (
|
||||
}
|
||||
|
||||
// GetNodes returns a list of registered servers
|
||||
func (s *PresenceService) GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
func (s *PresenceService) GetNodes(ctx context.Context, namespace string) ([]types.Server, error) {
|
||||
if namespace == "" {
|
||||
return nil, trace.BadParameter("missing namespace value")
|
||||
}
|
||||
@@ -223,9 +223,11 @@ func (s *PresenceService) GetNodes(ctx context.Context, namespace string, opts .
|
||||
server, err := services.UnmarshalServer(
|
||||
item.Value,
|
||||
types.KindNode,
|
||||
services.AddOptions(opts,
|
||||
[]services.MarshalOption{
|
||||
services.WithResourceID(item.ID),
|
||||
services.WithExpires(item.Expires))...)
|
||||
services.WithExpires(item.Expires),
|
||||
}...,
|
||||
)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
@@ -235,59 +237,6 @@ func (s *PresenceService) GetNodes(ctx context.Context, namespace string, opts .
|
||||
return servers, nil
|
||||
}
|
||||
|
||||
// ListNodes returns a paginated list of registered servers.
|
||||
// StartKey is a resource name, which is the suffix of its key.
|
||||
//
|
||||
// DELETE IN 10.0.0 in favor of ListResources.
|
||||
func (s *PresenceService) ListNodes(ctx context.Context, req proto.ListNodesRequest) (page []types.Server, nextKey string, err error) {
|
||||
// NOTE: changes to the outward behavior of this method may require updating cache.Cache.ListNodes, since that method
|
||||
// emulates this one but relies on a different implementation internally.
|
||||
if req.Namespace == "" {
|
||||
return nil, "", trace.BadParameter("missing namespace value")
|
||||
}
|
||||
limit := int(req.Limit)
|
||||
if limit <= 0 {
|
||||
return nil, "", trace.BadParameter("nonpositive limit value")
|
||||
}
|
||||
|
||||
// Get all items in the bucket within the given range.
|
||||
rangeStart := backend.Key(nodesPrefix, req.Namespace, req.StartKey)
|
||||
keyPrefix := backend.Key(nodesPrefix, req.Namespace)
|
||||
rangeEnd := backend.RangeEnd(keyPrefix)
|
||||
|
||||
var servers []types.Server
|
||||
err = backend.IterateRange(ctx, s.Backend, rangeStart, rangeEnd, limit, func(items []backend.Item) (stop bool, err error) {
|
||||
for _, item := range items {
|
||||
if len(servers) == limit {
|
||||
break
|
||||
}
|
||||
server, err := services.UnmarshalServer(
|
||||
item.Value,
|
||||
types.KindNode,
|
||||
services.WithResourceID(item.ID),
|
||||
services.WithExpires(item.Expires),
|
||||
)
|
||||
if err != nil {
|
||||
return false, trace.Wrap(err)
|
||||
}
|
||||
if server.MatchAgainst(req.Labels) {
|
||||
servers = append(servers, server)
|
||||
}
|
||||
}
|
||||
return len(servers) == limit, nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, "", trace.Wrap(err)
|
||||
}
|
||||
|
||||
// If a full page was filled, set nextKey using the last node.
|
||||
if len(servers) == limit {
|
||||
nextKey = backend.NextPaginationKey(servers[len(servers)-1])
|
||||
}
|
||||
|
||||
return servers, nextKey, nil
|
||||
}
|
||||
|
||||
// UpsertNode registers node presence, permanently if TTL is 0 or for the
|
||||
// specified duration with second resolution if it's >= 1 second.
|
||||
func (s *PresenceService) UpsertNode(ctx context.Context, server types.Server) (*types.KeepAlive, error) {
|
||||
@@ -333,43 +282,6 @@ func (s *PresenceService) KeepAliveNode(ctx context.Context, h types.KeepAlive)
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
|
||||
// UpsertNodes is used for bulk insertion of nodes.
|
||||
func (s *PresenceService) UpsertNodes(namespace string, servers []types.Server) error {
|
||||
batch, ok := s.Backend.(backend.Batch)
|
||||
if !ok {
|
||||
return trace.BadParameter("backend does not support batch interface")
|
||||
}
|
||||
if namespace == "" {
|
||||
return trace.BadParameter("missing node namespace")
|
||||
}
|
||||
|
||||
start := time.Now()
|
||||
|
||||
items := make([]backend.Item, len(servers))
|
||||
for i, server := range servers {
|
||||
value, err := services.MarshalServer(server)
|
||||
if err != nil {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
|
||||
items[i] = backend.Item{
|
||||
Key: backend.Key(nodesPrefix, server.GetNamespace(), server.GetName()),
|
||||
Value: value,
|
||||
Expires: server.Expiry(),
|
||||
ID: server.GetResourceID(),
|
||||
}
|
||||
}
|
||||
|
||||
err := batch.PutRange(context.TODO(), items)
|
||||
if err != nil {
|
||||
return trace.Wrap(err)
|
||||
}
|
||||
|
||||
s.log.Debugf("UpsertNodes(%v) in %v", len(servers), time.Since(start))
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetAuthServers returns a list of registered servers
|
||||
func (s *PresenceService) GetAuthServers() ([]types.Server, error) {
|
||||
return s.getServers(context.TODO(), types.KindAuthServer, authServersPrefix)
|
||||
|
||||
@@ -313,57 +313,6 @@ func TestNodeCRUD(t *testing.T) {
|
||||
|
||||
// Run NodeGetters in nested subtests to allow parallelization.
|
||||
t.Run("NodeGetters", func(t *testing.T) {
|
||||
t.Run("List Nodes", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
// list nodes one at a time, last page should be empty
|
||||
nodes, nextKey, err := presence.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: 1,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 1, len(nodes))
|
||||
require.Empty(t, cmp.Diff([]types.Server{node1}, nodes,
|
||||
cmpopts.IgnoreFields(types.Metadata{}, "ID")))
|
||||
require.EqualValues(t, backend.NextPaginationKey(node1), nextKey)
|
||||
|
||||
nodes, nextKey, err = presence.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: 1,
|
||||
StartKey: nextKey,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 1, len(nodes))
|
||||
require.Empty(t, cmp.Diff([]types.Server{node2}, nodes,
|
||||
cmpopts.IgnoreFields(types.Metadata{}, "ID")))
|
||||
require.EqualValues(t, backend.NextPaginationKey(node2), nextKey)
|
||||
|
||||
nodes, nextKey, err = presence.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: 1,
|
||||
StartKey: nextKey,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.EqualValues(t, 0, len(nodes))
|
||||
require.EqualValues(t, "", nextKey)
|
||||
|
||||
// ListNodes should fail if namespace isn't provided
|
||||
_, _, err = presence.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Limit: 1,
|
||||
})
|
||||
require.IsType(t, &trace.BadParameterError{}, err.(*trace.TraceErr).OrigError())
|
||||
|
||||
// ListNodes should fail if limit is nonpositive
|
||||
_, _, err = presence.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
})
|
||||
require.IsType(t, &trace.BadParameterError{}, err.(*trace.TraceErr).OrigError())
|
||||
|
||||
_, _, err = presence.ListNodes(ctx, proto.ListNodesRequest{
|
||||
Namespace: apidefaults.Namespace,
|
||||
Limit: -1,
|
||||
})
|
||||
require.IsType(t, &trace.BadParameterError{}, err.(*trace.TraceErr).OrigError())
|
||||
})
|
||||
t.Run("GetNodes", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
// Get all nodes, transparently handle limit exceeded errors
|
||||
|
||||
@@ -32,7 +32,7 @@ type ProxyGetter interface {
|
||||
// NodesGetter is a service that gets nodes.
|
||||
type NodesGetter interface {
|
||||
// GetNodes returns a list of registered servers.
|
||||
GetNodes(ctx context.Context, namespace string, opts ...MarshalOption) ([]types.Server, error)
|
||||
GetNodes(ctx context.Context, namespace string) ([]types.Server, error)
|
||||
}
|
||||
|
||||
// Presence records and reports the presence of all components
|
||||
@@ -43,8 +43,6 @@ type Presence interface {
|
||||
|
||||
// GetNode returns a node by name and namespace.
|
||||
GetNode(ctx context.Context, namespace, name string) (types.Server, error)
|
||||
// ListNodes returns a paginated list of registered servers.
|
||||
ListNodes(ctx context.Context, req proto.ListNodesRequest) (nodes []types.Server, nextKey string, err error)
|
||||
|
||||
// NodesGetter gets nodes
|
||||
NodesGetter
|
||||
@@ -59,9 +57,6 @@ type Presence interface {
|
||||
// specified duration with second resolution if it's >= 1 second.
|
||||
UpsertNode(ctx context.Context, server types.Server) (*types.KeepAlive, error)
|
||||
|
||||
// UpsertNodes bulk inserts nodes.
|
||||
UpsertNodes(namespace string, servers []types.Server) error
|
||||
|
||||
// DELETE IN: 5.1.0
|
||||
//
|
||||
// This logic has been moved to KeepAliveServer.
|
||||
|
||||
@@ -19,6 +19,7 @@ package clusters
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/gravitational/teleport/api/client/proto"
|
||||
"github.com/gravitational/teleport/api/defaults"
|
||||
"github.com/gravitational/teleport/api/types"
|
||||
"github.com/gravitational/teleport/lib/teleterm/api/uri"
|
||||
@@ -42,7 +43,9 @@ func (c *Cluster) GetServers(ctx context.Context) ([]Server, error) {
|
||||
}
|
||||
defer proxyClient.Close()
|
||||
|
||||
clusterServers, err := proxyClient.FindServersByLabels(ctx, defaults.Namespace, nil)
|
||||
clusterServers, err := proxyClient.FindNodesByFilters(ctx, proto.ListResourcesRequest{
|
||||
Namespace: defaults.Namespace,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
@@ -3349,7 +3349,7 @@ type authProviderMock struct {
|
||||
server types.ServerV2
|
||||
}
|
||||
|
||||
func (mock authProviderMock) GetNodes(ctx context.Context, n string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
func (mock authProviderMock) GetNodes(ctx context.Context, n string) ([]types.Server, error) {
|
||||
return []types.Server{&mock.server}, nil
|
||||
}
|
||||
|
||||
|
||||
+1
-2
@@ -42,7 +42,6 @@ import (
|
||||
"github.com/gravitational/teleport/lib/client"
|
||||
"github.com/gravitational/teleport/lib/defaults"
|
||||
"github.com/gravitational/teleport/lib/events"
|
||||
"github.com/gravitational/teleport/lib/services"
|
||||
"github.com/gravitational/teleport/lib/session"
|
||||
"github.com/gravitational/teleport/lib/sshutils"
|
||||
"github.com/gravitational/teleport/lib/utils"
|
||||
@@ -83,7 +82,7 @@ type TerminalRequest struct {
|
||||
|
||||
// AuthProvider is a subset of the full Auth API.
|
||||
type AuthProvider interface {
|
||||
GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error)
|
||||
GetNodes(ctx context.Context, namespace string) ([]types.Server, error)
|
||||
GetSessionEvents(namespace string, sid session.ID, after int, includePrintEvents bool) ([]events.EventFields, error)
|
||||
}
|
||||
|
||||
|
||||
@@ -92,7 +92,7 @@ func GetClusterDetails(ctx context.Context, site reversetunnel.RemoteSite, opts
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
nodes, err := clt.GetNodes(ctx, apidefaults.Namespace, opts...)
|
||||
nodes, err := clt.GetNodes(ctx, apidefaults.Namespace)
|
||||
if err != nil {
|
||||
return nil, trace.Wrap(err)
|
||||
}
|
||||
|
||||
@@ -182,8 +182,8 @@ type mockAccessPoint struct {
|
||||
presence *local.PresenceService
|
||||
}
|
||||
|
||||
func (m *mockAccessPoint) GetNodes(ctx context.Context, namespace string, opts ...services.MarshalOption) ([]types.Server, error) {
|
||||
return m.presence.GetNodes(ctx, namespace, opts...)
|
||||
func (m *mockAccessPoint) GetNodes(ctx context.Context, namespace string) ([]types.Server, error) {
|
||||
return m.presence.GetNodes(ctx, namespace)
|
||||
}
|
||||
|
||||
func (m *mockAccessPoint) GetProxies() ([]types.Server, error) {
|
||||
|
||||
Reference in New Issue
Block a user