From 87331a926563150d027499bba40e90d4a8952bb9 Mon Sep 17 00:00:00 2001 From: Pavel Sviderski Date: Fri, 11 Apr 2025 21:26:40 +1000 Subject: [PATCH] chore: add filter to ListMachines --- cmd/uncloud/machine/add.go | 2 +- cmd/uncloud/machine/ls.go | 2 +- cmd/uncloud/service/inspect.go | 2 +- cmd/uncloud/service/rm.go | 4 +- cmd/uncloud/volume/create.go | 2 +- cmd/uncloud/volume/ls.go | 11 +----- cmd/uncloud/volume/rm.go | 12 +----- internal/cli/cli.go | 2 +- internal/cli/flags.go | 23 ++++++++++++ internal/ucind/cluster.go | 2 +- pkg/api/client.go | 12 +++++- pkg/api/machine.go | 8 ++++ pkg/client/cluster.go | 33 ----------------- pkg/client/deploy/strategy.go | 4 +- pkg/client/machine.go | 67 ++++++++++++++++++++++++++++++++++ pkg/client/resolver.go | 2 +- pkg/client/service.go | 6 +-- pkg/client/volume.go | 2 +- test/e2e/cluster_test.go | 9 +++-- 19 files changed, 129 insertions(+), 76 deletions(-) create mode 100644 internal/cli/flags.go delete mode 100644 pkg/client/cluster.go create mode 100644 pkg/client/machine.go diff --git a/cmd/uncloud/machine/add.go b/cmd/uncloud/machine/add.go index 7eb335c4..e6d7415d 100644 --- a/cmd/uncloud/machine/add.go +++ b/cmd/uncloud/machine/add.go @@ -163,7 +163,7 @@ func waitClusterInitialised(ctx context.Context, client *client.Client) error { ), ctx) check := func() error { - _, err := client.ListMachines(ctx) + _, err := client.ListMachines(ctx, nil) if err == nil { return nil } diff --git a/cmd/uncloud/machine/ls.go b/cmd/uncloud/machine/ls.go index 8d24afbf..9504ca72 100644 --- a/cmd/uncloud/machine/ls.go +++ b/cmd/uncloud/machine/ls.go @@ -38,7 +38,7 @@ func list(ctx context.Context, uncli *cli.CLI, clusterName string) error { } defer client.Close() - machines, err := client.ListMachines(ctx) + machines, err := client.ListMachines(ctx, nil) if err != nil { return fmt.Errorf("list machines: %w", err) } diff --git a/cmd/uncloud/service/inspect.go b/cmd/uncloud/service/inspect.go index f0706db4..2b39d093 100644 --- a/cmd/uncloud/service/inspect.go +++ b/cmd/uncloud/service/inspect.go @@ -50,7 +50,7 @@ func inspect(ctx context.Context, uncli *cli.CLI, opts inspectOptions) error { return fmt.Errorf("inspect service: %w", err) } - machines, err := client.ListMachines(ctx) + machines, err := client.ListMachines(ctx, nil) if err != nil { return fmt.Errorf("list machines: %w", err) } diff --git a/cmd/uncloud/service/rm.go b/cmd/uncloud/service/rm.go index cbbbefcd..5609ab60 100644 --- a/cmd/uncloud/service/rm.go +++ b/cmd/uncloud/service/rm.go @@ -3,9 +3,7 @@ package service import ( "context" "fmt" - "os" - "github.com/docker/cli/cli/streams" "github.com/docker/compose/v2/pkg/progress" "github.com/psviderski/uncloud/internal/cli" "github.com/spf13/cobra" @@ -49,7 +47,7 @@ func rm(ctx context.Context, uncli *cli.CLI, opts rmOptions) error { return fmt.Errorf("remove service '%s': %w", s, err) } return nil - }, streams.NewOut(os.Stdout), "Removing service "+s) + }, uncli.ProgressOut(), "Removing service "+s) } return nil diff --git a/cmd/uncloud/volume/create.go b/cmd/uncloud/volume/create.go index 003dafb1..2fcc15b5 100644 --- a/cmd/uncloud/volume/create.go +++ b/cmd/uncloud/volume/create.go @@ -85,7 +85,7 @@ func create(ctx context.Context, uncli *cli.CLI, name string, opts createOptions // List machines and filter by the specified machine name or ID. // If no machine is specified, prompt the user to select one. - machines, err := client.ListMachines(ctx) + machines, err := client.ListMachines(ctx, nil) if err != nil { return fmt.Errorf("list machines: %w", err) } diff --git a/cmd/uncloud/volume/ls.go b/cmd/uncloud/volume/ls.go index 71c4d5e6..71aaf11e 100644 --- a/cmd/uncloud/volume/ls.go +++ b/cmd/uncloud/volume/ls.go @@ -54,16 +54,7 @@ func list(ctx context.Context, uncli *cli.CLI, opts listOptions) error { // Apply machine filter if specified. var filter *api.VolumeFilter if len(opts.machines) > 0 { - // Expand comma-separated machine names. - var machines []string - for _, m := range opts.machines { - for _, nameOrID := range strings.Split(m, ",") { - if nameOrID = strings.TrimSpace(nameOrID); nameOrID != "" { - machines = append(machines, nameOrID) - } - } - } - + machines := cli.ExpandCommaSeparatedValues(opts.machines) filter = &api.VolumeFilter{ Machines: machines, } diff --git a/cmd/uncloud/volume/rm.go b/cmd/uncloud/volume/rm.go index ed4fa6fe..148769f9 100644 --- a/cmd/uncloud/volume/rm.go +++ b/cmd/uncloud/volume/rm.go @@ -4,7 +4,6 @@ import ( "context" "errors" "fmt" - "strings" "github.com/psviderski/uncloud/internal/cli" "github.com/psviderski/uncloud/pkg/api" @@ -60,16 +59,7 @@ func remove(ctx context.Context, uncli *cli.CLI, names []string, opts removeOpti } if len(opts.machines) > 0 { - // Expand comma-separated machine names. - var machines []string - for _, m := range opts.machines { - for _, nameOrID := range strings.Split(m, ",") { - if nameOrID = strings.TrimSpace(nameOrID); nameOrID != "" { - machines = append(machines, nameOrID) - } - } - } - + machines := cli.ExpandCommaSeparatedValues(opts.machines) filter.Machines = machines } diff --git a/internal/cli/cli.go b/internal/cli/cli.go index 07b8315a..6a144b7b 100644 --- a/internal/cli/cli.go +++ b/internal/cli/cli.go @@ -305,7 +305,7 @@ func (cli *CLI) AddMachine( } // List other machines in the cluster to include them in the join request. - machines, err := c.ListMachines(ctx) + machines, err := c.ListMachines(ctx, nil) if err != nil { return nil, fmt.Errorf("list cluster machines: %w", err) } diff --git a/internal/cli/flags.go b/internal/cli/flags.go new file mode 100644 index 00000000..289bfa73 --- /dev/null +++ b/internal/cli/flags.go @@ -0,0 +1,23 @@ +package cli + +import ( + "strings" +) + +// ExpandCommaSeparatedValues takes a slice of strings and expands any comma-separated values into individual elements. +func ExpandCommaSeparatedValues(values []string) []string { + if len(values) == 0 { + return nil + } + + var expanded []string + for _, value := range values { + for _, v := range strings.Split(value, ",") { + if v = strings.TrimSpace(v); v != "" { + expanded = append(expanded, v) + } + } + } + + return expanded +} diff --git a/internal/ucind/cluster.go b/internal/ucind/cluster.go index 8246fd5c..8b409ffb 100644 --- a/internal/ucind/cluster.go +++ b/internal/ucind/cluster.go @@ -239,7 +239,7 @@ func (p *Provisioner) WaitClusterReady(ctx context.Context, c Cluster, timeout t ), ctx) checkMachinesUp := func() error { - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return fmt.Errorf("list machines: %w", err) } diff --git a/pkg/api/client.go b/pkg/api/client.go index b7c85b0f..524870c1 100644 --- a/pkg/api/client.go +++ b/pkg/api/client.go @@ -6,6 +6,7 @@ import ( "slices" "github.com/docker/docker/api/types/container" + "github.com/docker/docker/api/types/volume" "github.com/psviderski/uncloud/internal/machine/api/pb" "google.golang.org/grpc/metadata" ) @@ -16,6 +17,7 @@ type Client interface { ImageClient MachineClient ServiceClient + VolumeClient } type ContainerClient interface { @@ -39,19 +41,25 @@ type ImageClient interface { type MachineClient interface { InspectMachine(ctx context.Context, id string) (*pb.MachineMember, error) - ListMachines(ctx context.Context) (MachineMembersList, error) + ListMachines(ctx context.Context, filter *MachineFilter) (MachineMembersList, error) } type ServiceClient interface { InspectService(ctx context.Context, id string) (Service, error) } +type VolumeClient interface { + CreateVolume(ctx context.Context, machineNameOrID string, opts volume.CreateOptions) (MachineVolume, error) + ListVolumes(ctx context.Context, filter *VolumeFilter) ([]MachineVolume, error) + RemoveVolume(ctx context.Context, machineNameOrID, volumeName string, force bool) error +} + // ProxyMachinesContext returns a new context that proxies gRPC requests to the specified machines. // If namesOrIDs is nil, all machines are included. func ProxyMachinesContext( ctx context.Context, cli MachineClient, namesOrIDs []string, ) (context.Context, MachineMembersList, error) { - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return nil, nil, fmt.Errorf("list machines: %w", err) } diff --git a/pkg/api/machine.go b/pkg/api/machine.go index 885955c4..a81f2872 100644 --- a/pkg/api/machine.go +++ b/pkg/api/machine.go @@ -2,6 +2,14 @@ package api import "github.com/psviderski/uncloud/internal/machine/api/pb" +// MachineFilter defines criteria to filter machines in ListMachines. +type MachineFilter struct { + // Available filters machines that are not DOWN. + Available bool + // NamesOrIDs filters machines by their names or IDs. + NamesOrIDs []string +} + type MachineMembersList []*pb.MachineMember func (m MachineMembersList) FindByManagementIP(ip string) *pb.MachineMember { diff --git a/pkg/client/cluster.go b/pkg/client/cluster.go deleted file mode 100644 index 70bdfbf9..00000000 --- a/pkg/client/cluster.go +++ /dev/null @@ -1,33 +0,0 @@ -package client - -import ( - "context" - - "github.com/psviderski/uncloud/internal/machine/api/pb" - "github.com/psviderski/uncloud/pkg/api" - "google.golang.org/protobuf/types/known/emptypb" -) - -func (cli *Client) InspectMachine(ctx context.Context, nameOrID string) (*pb.MachineMember, error) { - machines, err := cli.ListMachines(ctx) - if err != nil { - return nil, err - } - - for _, m := range machines { - if m.Machine.Id == nameOrID || m.Machine.Name == nameOrID { - return m, nil - } - } - - return nil, api.ErrNotFound -} - -// ListMachines returns a list of all machines registered in the cluster. -func (cli *Client) ListMachines(ctx context.Context) (api.MachineMembersList, error) { - resp, err := cli.ClusterClient.ListMachines(ctx, &emptypb.Empty{}) - if err != nil { - return nil, err - } - return resp.Machines, nil -} diff --git a/pkg/client/deploy/strategy.go b/pkg/client/deploy/strategy.go index 208c24f9..dbff1479 100644 --- a/pkg/client/deploy/strategy.go +++ b/pkg/client/deploy/strategy.go @@ -58,7 +58,7 @@ func (s *RollingStrategy) planReplicated( return plan, err } - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return plan, fmt.Errorf("list machines: %w", err) } @@ -227,7 +227,7 @@ func (s *RollingStrategy) planGlobal( } } - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return plan, fmt.Errorf("list machines: %w", err) } diff --git a/pkg/client/machine.go b/pkg/client/machine.go new file mode 100644 index 00000000..db4638b4 --- /dev/null +++ b/pkg/client/machine.go @@ -0,0 +1,67 @@ +package client + +import ( + "context" + "slices" + + "github.com/psviderski/uncloud/internal/machine/api/pb" + "github.com/psviderski/uncloud/pkg/api" + "google.golang.org/protobuf/types/known/emptypb" +) + +func (cli *Client) InspectMachine(ctx context.Context, nameOrID string) (*pb.MachineMember, error) { + machines, err := cli.ListMachines(ctx, nil) + if err != nil { + return nil, err + } + + for _, m := range machines { + if m.Machine.Id == nameOrID || m.Machine.Name == nameOrID { + return m, nil + } + } + + return nil, api.ErrNotFound +} + +// ListMachines returns a list of all machines registered in the cluster that match the filter. +func (cli *Client) ListMachines(ctx context.Context, filter *api.MachineFilter) (api.MachineMembersList, error) { + resp, err := cli.ClusterClient.ListMachines(ctx, &emptypb.Empty{}) + if err != nil { + return nil, err + } + + machines := resp.Machines + + if filter != nil { + var matchedMachines api.MachineMembersList + for _, m := range machines { + if MachineMatchesFilter(m, filter) { + matchedMachines = append(matchedMachines, m) + } + } + machines = matchedMachines + } + + return machines, nil +} + +func MachineMatchesFilter(machine *pb.MachineMember, filter *api.MachineFilter) bool { + if filter == nil { + return true + } + + if filter.Available && machine.State == pb.MachineMember_DOWN { + return false + } + + if len(filter.NamesOrIDs) > 0 { + if !slices.ContainsFunc(filter.NamesOrIDs, func(nameOrID string) bool { + return machine.Machine.Id == nameOrID || machine.Machine.Name == nameOrID + }) { + return false + } + } + + return true +} diff --git a/pkg/client/resolver.go b/pkg/client/resolver.go index e1135269..26cb8385 100644 --- a/pkg/client/resolver.go +++ b/pkg/client/resolver.go @@ -37,7 +37,7 @@ func (r *MapNameResolver) ContainerName(containerID string) string { // ServiceOperationNameResolver returns a machine and container name resolver for a service that can be used to format // deployment operations. func (cli *Client) ServiceOperationNameResolver(ctx context.Context, svc api.Service) (*MapNameResolver, error) { - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return nil, fmt.Errorf("list machines: %w", err) } diff --git a/pkg/client/service.go b/pkg/client/service.go index 30cc728c..babfbc19 100644 --- a/pkg/client/service.go +++ b/pkg/client/service.go @@ -69,7 +69,7 @@ func (cli *Client) RunService( func (cli *Client) InspectService(ctx context.Context, nameOrID string) (api.Service, error) { var svc api.Service - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return svc, fmt.Errorf("list machines: %w", err) } @@ -204,7 +204,7 @@ func (cli *Client) RemoveService(ctx context.Context, id string) error { return err } - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return fmt.Errorf("list machines: %w", err) } @@ -251,7 +251,7 @@ func (cli *Client) RemoveService(ctx context.Context, id string) error { // ListServices returns a list of all services and their containers. func (cli *Client) ListServices(ctx context.Context) ([]api.Service, error) { - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { return nil, fmt.Errorf("list machines: %w", err) } diff --git a/pkg/client/volume.go b/pkg/client/volume.go index a8204ece..1978d9c7 100644 --- a/pkg/client/volume.go +++ b/pkg/client/volume.go @@ -47,7 +47,7 @@ func (cli *Client) CreateVolume( return resp, nil } -// ListVolumes returns a list of all volumes on the cluster machines. +// ListVolumes returns a list of all volumes on the cluster machines that match the filter. func (cli *Client) ListVolumes(ctx context.Context, filter *api.VolumeFilter) ([]api.MachineVolume, error) { // Broadcast the volume list request to the specified machines in the filter or all machines if filter is nil. var proxyMachines []string diff --git a/test/e2e/cluster_test.go b/test/e2e/cluster_test.go index f7025c56..5bc11cc5 100644 --- a/test/e2e/cluster_test.go +++ b/test/e2e/cluster_test.go @@ -3,6 +3,10 @@ package e2e import ( "context" "errors" + "os" + "testing" + "time" + dockerclient "github.com/docker/docker/client" "github.com/psviderski/uncloud/internal/machine/api/pb" "github.com/psviderski/uncloud/internal/ucind" @@ -11,9 +15,6 @@ import ( "github.com/stretchr/testify/require" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" - "os" - "testing" - "time" ) func createTestCluster( @@ -79,7 +80,7 @@ func TestClusterLifecycle(t *testing.T) { for i, cli := range clients { // Wait for the machine to reconcile the cluster store. require.Eventually(t, func() bool { - machines, err := cli.ListMachines(ctx) + machines, err := cli.ListMachines(ctx, nil) if err != nil { // FailedPrecondition "cluster is not initialised" is expected until the store is reconciled. if s, ok := status.FromError(err); ok {