Files

126 lines
3.6 KiB
Go

package client
import (
"context"
"fmt"
"github.com/containerd/errdefs"
"github.com/docker/compose/v2/pkg/progress"
"github.com/docker/docker/api/types/volume"
cliprogress "github.com/psviderski/uncloud/internal/cli/progress"
"github.com/psviderski/uncloud/internal/cli/tui"
"github.com/psviderski/uncloud/pkg/api"
)
// CreateVolume creates a new volume on the specified machine.
func (cli *Client) CreateVolume(
ctx context.Context, machineNameOrID string, opts volume.CreateOptions,
) (api.MachineVolume, error) {
var resp api.MachineVolume
if opts.Name == "" {
return resp, fmt.Errorf("volume name is required (anonymous volumes are not supported)")
}
machine, err := cli.InspectMachine(ctx, machineNameOrID)
if err != nil {
return resp, fmt.Errorf("inspect machine '%s': %w", machineNameOrID, err)
}
// Proxy Docker gRPC requests to the selected machine.
ctx = cli.ProxySingleMachineContext(ctx, machine.Machine.Id)
pw := progress.ContextWriter(ctx)
eventID := cliprogress.VolumeEventID(opts.Name, machine.Machine.Name)
pw.Event(progress.CreatingEvent(eventID))
vol, err := cli.Docker.CreateVolume(ctx, opts)
if err != nil {
return resp, err
}
resp = api.MachineVolume{
MachineID: machine.Machine.Id,
MachineName: machine.Machine.Name,
Volume: vol,
}
pw.Event(progress.CreatedEvent(eventID))
return resp, nil
}
// 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
if filter != nil {
proxyMachines = filter.Machines
}
listCtx := cli.ProxyMachinesContext(ctx, proxyMachines)
machineVolumes, err := cli.Docker.ListVolumes(listCtx, volume.ListOptions{})
if err != nil {
return nil, err
}
var volumes []api.MachineVolume
// Process responses from all machines.
for _, mv := range machineVolumes {
if mv.Metadata == nil {
tui.PrintWarning("metadata is missing in response from unknown server")
continue
}
if mv.Metadata.Error != "" {
// TODO: return failed machines in the response.
tui.PrintWarning(fmt.Sprintf("failed to list volumes on machine '%s': %s", mv.Metadata.MachineName,
mv.Metadata.Error))
continue
}
for _, vol := range mv.Response.Volumes {
volumes = append(volumes, api.MachineVolume{
MachineID: mv.Metadata.MachineId,
MachineName: mv.Metadata.MachineName,
Volume: *vol,
})
}
}
// 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
}
// RemoveVolume removes a volume from the specified machine.
func (cli *Client) RemoveVolume(ctx context.Context, machineNameOrID, volumeName string, force bool) error {
machine, err := cli.InspectMachine(ctx, machineNameOrID)
if err != nil {
return fmt.Errorf("inspect machine '%s': %w", machineNameOrID, err)
}
// Proxy Docker gRPC requests to the selected machine.
ctx = cli.ProxySingleMachineContext(ctx, machine.Machine.Id)
pw := progress.ContextWriter(ctx)
eventID := cliprogress.VolumeEventID(volumeName, machine.Machine.Name)
pw.Event(progress.RemovingEvent(eventID))
if err = cli.Docker.RemoveVolume(ctx, volumeName, force); err != nil {
if errdefs.IsNotFound(err) {
return api.ErrNotFound
}
return err
}
pw.Event(progress.RemovedEvent(eventID))
return nil
}