diff --git a/pkg/api/client.go b/pkg/api/client.go index 66a2ea82..b7c85b0f 100644 --- a/pkg/api/client.go +++ b/pkg/api/client.go @@ -48,20 +48,24 @@ type ServiceClient interface { // 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, error) { +func ProxyMachinesContext( + ctx context.Context, cli MachineClient, namesOrIDs []string, +) (context.Context, MachineMembersList, error) { machines, err := cli.ListMachines(ctx) if err != nil { - return nil, fmt.Errorf("list machines: %w", err) + return nil, nil, fmt.Errorf("list machines: %w", err) } + var proxiedMachines MachineMembersList md := metadata.New(nil) for _, m := range machines { - if namesOrIDs == nil || + if len(namesOrIDs) == 0 || slices.Contains(namesOrIDs, m.Machine.Name) || slices.Contains(namesOrIDs, m.Machine.Id) { + proxiedMachines = append(proxiedMachines, m) machineIP, _ := m.Machine.Network.ManagementIp.ToAddr() md.Append("machines", machineIP.String()) } } - return metadata.NewOutgoingContext(ctx, md), nil + return metadata.NewOutgoingContext(ctx, md), proxiedMachines, nil } diff --git a/pkg/api/machine.go b/pkg/api/machine.go index 6978982b..885955c4 100644 --- a/pkg/api/machine.go +++ b/pkg/api/machine.go @@ -18,3 +18,13 @@ func (m MachineMembersList) FindByManagementIP(ip string) *pb.MachineMember { return nil } + +func (m MachineMembersList) FindByNameOrID(nameOrID string) *pb.MachineMember { + for _, machine := range m { + if machine.Machine.Id == nameOrID || machine.Machine.Name == nameOrID { + return machine + } + } + + return nil +} diff --git a/pkg/api/volume.go b/pkg/api/volume.go index 4e71ef0e..e82231bb 100644 --- a/pkg/api/volume.go +++ b/pkg/api/volume.go @@ -3,6 +3,7 @@ package api import ( "fmt" "reflect" + "slices" "sort" "strings" @@ -195,3 +196,47 @@ type MachineVolume struct { // Volume is the Docker volume model. Volume volume.Volume } + +// VolumeFilter defines criteria to filter volumes in ListVolumes. +type VolumeFilter struct { + // Driver filters volumes by storage driver name. + Driver string + // Labels filters volumes by label key-value pairs. Volumes must match all labels. + Labels map[string]string + // MachineIDs filters volumes to those on the specified machines (names or IDs). + Machines []string + // Names filters volumes by name. Volumes must match one of the names. + Names []string +} + +// MatchesFilter checks if a volume matches the specified filter criteria. +func (v *MachineVolume) MatchesFilter(filter *VolumeFilter) bool { + if filter == nil { + return true + } + // Filter by name. + if len(filter.Names) > 0 && !slices.Contains(filter.Names, v.Volume.Name) { + return false + } + // Filter by driver. + if filter.Driver != "" && v.Volume.Driver != filter.Driver { + return false + } + // Filter by labels. + for key, value := range filter.Labels { + labelValue, exists := v.Volume.Labels[key] + if !exists || labelValue != value { + return false + } + } + // Filter by machines. + if len(filter.Machines) > 0 { + if !slices.ContainsFunc(filter.Machines, func(nameOrID string) bool { + return v.MachineName == nameOrID || v.MachineID == nameOrID + }) { + return false + } + } + + return true +} diff --git a/pkg/client/deploy/resolver.go b/pkg/client/deploy/resolver.go index e3fd4e12..7b0c9a35 100644 --- a/pkg/client/deploy/resolver.go +++ b/pkg/client/deploy/resolver.go @@ -4,14 +4,15 @@ import ( "context" "errors" "fmt" + "strings" + "time" + "github.com/distribution/reference" "github.com/docker/docker/api/types" "github.com/opencontainers/go-digest" "github.com/psviderski/uncloud/internal/secret" "github.com/psviderski/uncloud/pkg/api" "google.golang.org/grpc/codes" - "strings" - "time" ) // ServiceSpecResolver transforms user-provided service specs into deployment-ready form. @@ -187,7 +188,7 @@ func (r *ImageDigestResolver) Resolve(image, policy string) (string, error) { // resolveAlways resolves the image to the image with the digest by querying the registry from all machines. func (r *ImageDigestResolver) resolveAlways(image string) (string, error) { // TODO: broadcast to a subset of machines in large clusters to avoid being rate-limited by the registry. - ctx, err := api.ProxyMachinesContext(r.Ctx, r.Client, nil) + ctx, _, err := api.ProxyMachinesContext(r.Ctx, r.Client, nil) if err != nil { return "", fmt.Errorf("create request context to broadcast to all machines: %w", err) } @@ -215,7 +216,7 @@ func (r *ImageDigestResolver) resolveAlways(image string) (string, error) { } func (r *ImageDigestResolver) resolveMissing(image string) (string, error) { - ctx, err := api.ProxyMachinesContext(r.Ctx, r.Client, nil) + ctx, _, err := api.ProxyMachinesContext(r.Ctx, r.Client, nil) if err != nil { return "", fmt.Errorf("create request context to broadcast to all machines: %w", err) } diff --git a/pkg/client/volume.go b/pkg/client/volume.go index 3f3af68d..a8204ece 100644 --- a/pkg/client/volume.go +++ b/pkg/client/volume.go @@ -34,7 +34,7 @@ func (cli *Client) CreateVolume( vol, err := cli.Docker.CreateVolume(ctx, opts) if err != nil { - return resp, fmt.Errorf("create volume on machine '%s': %w", machine.Machine.Name, err) + return resp, err } resp = api.MachineVolume{ @@ -48,21 +48,21 @@ func (cli *Client) CreateVolume( } // ListVolumes returns a list of all volumes on the cluster machines. -func (cli *Client) ListVolumes(ctx context.Context) ([]api.MachineVolume, error) { - machines, err := cli.ListMachines(ctx) - if err != nil { - return nil, fmt.Errorf("list machines: %w", err) +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 + if filter != nil { + proxyMachines = filter.Machines } - // Broadcast the volume list request to all machines. - listCtx, err := api.ProxyMachinesContext(ctx, cli, nil) + listCtx, machines, err := api.ProxyMachinesContext(ctx, cli, proxyMachines) if err != nil { return nil, fmt.Errorf("create request context to broadcast to all machines: %w", err) } machineVolumes, err := cli.Docker.ListVolumes(listCtx, volume.ListOptions{}) if err != nil { - return nil, fmt.Errorf("list volumes: %w", err) + return nil, err } var volumes []api.MachineVolume @@ -95,6 +95,17 @@ func (cli *Client) ListVolumes(ctx context.Context) ([]api.MachineVolume, error) } } + // Filter volumes based on the provided filter criteria. + if filter != nil { + var filteredVolumes []api.MachineVolume + for _, vol := range volumes { + if vol.MatchesFilter(filter) { + filteredVolumes = append(filteredVolumes, vol) + } + } + volumes = filteredVolumes + } + return volumes, nil } @@ -115,8 +126,7 @@ func (cli *Client) RemoveVolume(ctx context.Context, machineNameOrID, volumeName if dockerclient.IsErrNotFound(err) { return api.ErrNotFound } - return fmt.Errorf("remove volume '%s' from machine '%s': %w", - volumeName, machine.Machine.Name, err) + return err } pw.Event(progress.RemovedEvent(eventID))