diff --git a/pkg/cloud/capabilities.go b/pkg/cloud/capabilities.go new file mode 100644 index 0000000..8ae081c --- /dev/null +++ b/pkg/cloud/capabilities.go @@ -0,0 +1,41 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package cloud + +import ( + "context" + + "k8s.io/klog/v2" +) + +func (c *client) GetMaxCustomDiskSizeGB(ctx context.Context) (int64, error) { + logger := klog.FromContext(ctx) + logger.V(2).Info("CloudStack API call", "command", "ListCapabilities") + + resp, err := c.Configuration.ListCapabilities(c.Configuration.NewListCapabilitiesParams()) + if err != nil { + return 0, err + } + if resp.Capabilities == nil { + return 0, nil + } + + return resp.Capabilities.Customdiskofferingmaxsize, nil +} diff --git a/pkg/cloud/cloud.go b/pkg/cloud/cloud.go index bed6aa8..c3c0f07 100644 --- a/pkg/cloud/cloud.go +++ b/pkg/cloud/cloud.go @@ -46,6 +46,10 @@ type Interface interface { DetachVolume(ctx context.Context, volumeID string) error ExpandVolume(ctx context.Context, volumeID string, newSizeInGB int64) error + // GetMaxCustomDiskSizeGB returns CloudStack's maximum size for volumes with + // a custom disk offering (0 if CloudStack reports no limit). + GetMaxCustomDiskSizeGB(ctx context.Context) (int64, error) + CreateVolumeFromSnapshot(ctx context.Context, zoneID, name, projectID, snapshotID string, sizeInGB int64) (*Volume, error) GetSnapshotByID(ctx context.Context, snapshotID string) (*Snapshot, error) GetSnapshotByName(ctx context.Context, name string) (*Snapshot, error) @@ -69,6 +73,9 @@ type Volume struct { VirtualMachineID string DeviceID string + + // State is the CloudStack volume state, e.g. Allocated, Ready, Destroy. + State string } // Snapshot represents a CloudStack snapshot. @@ -83,6 +90,9 @@ type Snapshot struct { VolumeID string CreatedAt string + + // State is the CloudStack snapshot state, e.g. Creating, BackingUp, BackedUp. + State string } // VM represents a CloudStack Virtual Machine. diff --git a/pkg/cloud/fake/fake.go b/pkg/cloud/fake/fake.go index 84b38b2..34d4ced 100644 --- a/pkg/cloud/fake/fake.go +++ b/pkg/cloud/fake/fake.go @@ -158,6 +158,10 @@ func (f *fakeConnector) ExpandVolume(_ context.Context, volumeID string, newSize return cloud.ErrNotFound } +func (f *fakeConnector) GetMaxCustomDiskSizeGB(_ context.Context) (int64, error) { + return 0, nil +} + func (f *fakeConnector) CreateVolumeFromSnapshot(_ context.Context, zoneID, name, _, _ string, sizeInGB int64) (*cloud.Volume, error) { vol := &cloud.Volume{ ID: "fake-vol-from-snap-" + name, diff --git a/pkg/cloud/snapshots.go b/pkg/cloud/snapshots.go index fae632a..3a972d1 100644 --- a/pkg/cloud/snapshots.go +++ b/pkg/cloud/snapshots.go @@ -21,7 +21,6 @@ package cloud import ( "context" - "strings" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" @@ -46,10 +45,13 @@ func (c *client) GetSnapshotByID(ctx context.Context, snapshotID string) (*Snaps return &Snapshot{ ID: snapshot.Id, Name: snapshot.Name, + State: snapshot.State, + Size: snapshot.Virtualsize, DomainID: snapshot.Domainid, ProjectID: snapshot.Projectid, ZoneID: snapshot.Zoneid, VolumeID: snapshot.Volumeid, + CreatedAt: snapshot.Created, }, nil } @@ -70,6 +72,7 @@ func (c *client) CreateSnapshot(ctx context.Context, volumeID, name string) (*Sn return &Snapshot{ ID: snapshot.Id, Name: snapshot.Name, + State: snapshot.State, Size: snapshot.Virtualsize, DomainID: snapshot.Domainid, ProjectID: snapshot.Projectid, @@ -82,11 +85,10 @@ func (c *client) CreateSnapshot(ctx context.Context, volumeID, name string) (*Sn func (c *client) DeleteSnapshot(_ context.Context, snapshotID string) error { p := c.Snapshot.NewDeleteSnapshotParams(snapshotID) _, err := c.Snapshot.DeleteSnapshot(p) - if err != nil && strings.Contains(err.Error(), "4350") { - // CloudStack error InvalidParameterValueException - return ErrNotFound - } + // Errors are returned as they are. CloudStack reports an already deleted + // snapshot with several different messages, so the caller checks whether + // the snapshot still exists instead of matching error text. return err } @@ -107,6 +109,8 @@ func (c *client) GetSnapshotByName(ctx context.Context, name string) (*Snapshot, return &Snapshot{ ID: snapshot.Id, Name: snapshot.Name, + State: snapshot.State, + Size: snapshot.Virtualsize, DomainID: snapshot.Domainid, ProjectID: snapshot.Projectid, ZoneID: snapshot.Zoneid, @@ -148,6 +152,7 @@ func (c *client) ListSnapshots(ctx context.Context, volumeID, snapshotID string) s := &Snapshot{ ID: snapshot.Id, Name: snapshot.Name, + State: snapshot.State, Size: snapshot.Virtualsize, DomainID: snapshot.Domainid, ProjectID: snapshot.Projectid, diff --git a/pkg/cloud/volumes.go b/pkg/cloud/volumes.go index 3db152f..03db3f9 100644 --- a/pkg/cloud/volumes.go +++ b/pkg/cloud/volumes.go @@ -23,7 +23,6 @@ import ( "context" "fmt" "strconv" - "strings" "github.com/apache/cloudstack-go/v2/cloudstack" "k8s.io/klog/v2" @@ -42,6 +41,7 @@ func mapVolume(vol *cloudstack.Volume) *Volume { ZoneID: vol.Zoneid, VirtualMachineID: vol.Virtualmachineid, DeviceID: strconv.FormatInt(vol.Deviceid, 10), + State: vol.State, } } @@ -113,11 +113,11 @@ func (c *client) DeleteVolume(ctx context.Context, id string) error { "id": id, }) _, err := c.Volume.DeleteVolume(p) - if err != nil && strings.Contains(err.Error(), "4350") { - // CloudStack error InvalidParameterValueException - return ErrNotFound - } + // Errors are returned as they are. CloudStack uses the same error code for + // "no such volume" and for refusals such as "volume is attached", so the + // caller checks whether the volume still exists before treating a failed + // delete as done. return err } diff --git a/pkg/driver/controller.go b/pkg/driver/controller.go index 9ca1185..562e39e 100644 --- a/pkg/driver/controller.go +++ b/pkg/driver/controller.go @@ -69,7 +69,6 @@ func NewControllerServer(connector cloud.Interface) csi.ControllerServer { } } -//nolint:gocognit func (cs *controllerServer) CreateVolume(ctx context.Context, req *csi.CreateVolumeRequest) (*csi.CreateVolumeResponse, error) { logger := klog.FromContext(ctx) logger.V(6).Info("CreateVolume: called", "args", *req) @@ -104,34 +103,6 @@ func (cs *controllerServer) CreateVolume(ctx context.Context, req *csi.CreateVol } defer cs.volumeLocks.Release(name) - // Check if a volume with that name already exists. - vol, err := cs.connector.GetVolumeByName(ctx, name) - if err != nil { - if !errors.Is(err, cloud.ErrNotFound) { - // Error with CloudStack - return nil, status.Errorf(codes.Internal, "CloudStack error: %v", err) - } - } else { - // The volume exists. Check if it suits the request. - if ok, message := checkVolumeSuitable(vol, diskOfferingID, req.GetCapacityRange(), req.GetAccessibilityRequirements()); !ok { - return nil, status.Errorf(codes.AlreadyExists, "Volume %v already exists but does not satisfy request: %s", name, message) - } - // Existing volume is ok. - resp := &csi.CreateVolumeResponse{ - Volume: &csi.Volume{ - VolumeId: vol.ID, - CapacityBytes: vol.Size, - VolumeContext: req.GetParameters(), - // ContentSource: req.GetVolumeContentSource(), TODO: snapshot support. - AccessibleTopology: []*csi.Topology{ - Topology{ZoneID: vol.ZoneID}.ToCSI(), - }, - }, - } - - return resp, nil - } - // Check if this is a volume from snapshot var snapshotID string if src := req.GetVolumeContentSource(); src != nil { @@ -140,6 +111,18 @@ func (cs *controllerServer) CreateVolume(ctx context.Context, req *csi.CreateVol } } + // Check if a volume with that name already exists, e.g. because an earlier + // attempt timed out after CloudStack had started creating it. + vol, err := cs.connector.GetVolumeByName(ctx, name) + switch { + case err == nil && snapshotID != "": + return cs.adoptRestoredVolume(ctx, req, vol, snapshotID) + case err == nil: + return existingVolumeResponse(req, vol, diskOfferingID) + case !errors.Is(err, cloud.ErrNotFound): + return nil, status.Errorf(codes.Internal, "CloudStack error: %v", err) + } + // We have to create the volume. // Determine volume size using requested capacity range. @@ -148,45 +131,8 @@ func (cs *controllerServer) CreateVolume(ctx context.Context, req *csi.CreateVol return nil, status.Error(codes.InvalidArgument, err.Error()) } - // If creating from snapshot, get the snapshot size - var snapshotSizeGiB int64 if snapshotID != "" { - logger.Info("Creating volume from snapshot", "snapshotID", snapshotID) - // Call the cloud connector's CreateVolumeFromSnapshot if implemented - printVolumeAsJSON(req) - snapshot, err := cs.connector.GetSnapshotByID(ctx, snapshotID) - if errors.Is(err, cloud.ErrNotFound) { - return nil, status.Errorf(codes.NotFound, "Snapshot %v not found", snapshotID) - } else if err != nil { - // Error with CloudStack - return nil, status.Errorf(codes.Internal, "Error %v", err) - } - - logger.Info("PVC created with", "size", sizeInGB) - snapshotSizeGiB = util.RoundUpBytesToGB(snapshot.Size) - if snapshotSizeGiB > sizeInGB { - logger.Info("Snapshot size is greater than the request PVC, creating volume from snapshot of size", "snapshot size:", snapshotSizeGiB) - sizeInGB = snapshotSizeGiB - } - - volFromSnapshot, err := cs.connector.CreateVolumeFromSnapshot(ctx, snapshot.ZoneID, name, snapshot.ProjectID, snapshotID, sizeInGB) - if err != nil { - return nil, status.Errorf(codes.Internal, "Cannot create volume from snapshot %s: %v", snapshotID, err.Error()) - } - - resp := &csi.CreateVolumeResponse{ - Volume: &csi.Volume{ - VolumeId: volFromSnapshot.ID, - CapacityBytes: volFromSnapshot.Size, - VolumeContext: req.GetParameters(), - ContentSource: req.GetVolumeContentSource(), - AccessibleTopology: []*csi.Topology{ - Topology{ZoneID: volFromSnapshot.ZoneID}.ToCSI(), - }, - }, - } - - return resp, nil + return cs.createVolumeFromSnapshot(ctx, req, name, snapshotID, sizeInGB) } // Determine zone using topology constraints. @@ -252,6 +198,25 @@ func printVolumeAsJSON(vol *csi.CreateVolumeRequest) { klog.V(5).Infof("CreateVolumeRequest as JSON:\n%s", string(b)) } +// existingVolumeResponse returns an existing volume with the requested name if +// it suits the request. +func existingVolumeResponse(req *csi.CreateVolumeRequest, vol *cloud.Volume, diskOfferingID string) (*csi.CreateVolumeResponse, error) { + if ok, message := checkVolumeSuitable(vol, diskOfferingID, req.GetCapacityRange(), req.GetAccessibilityRequirements()); !ok { + return nil, status.Errorf(codes.AlreadyExists, "Volume %v already exists but does not satisfy request: %s", req.GetName(), message) + } + + return &csi.CreateVolumeResponse{ + Volume: &csi.Volume{ + VolumeId: vol.ID, + CapacityBytes: vol.Size, + VolumeContext: req.GetParameters(), + AccessibleTopology: []*csi.Topology{ + Topology{ZoneID: vol.ZoneID}.ToCSI(), + }, + }, + }, nil +} + func checkVolumeSuitable(vol *cloud.Volume, diskOfferingID string, capRange *csi.CapacityRange, topologyRequirement *csi.TopologyRequirement, ) (bool, string) { @@ -285,6 +250,147 @@ func checkVolumeSuitable(vol *cloud.Volume, return true, "" } +// isSnapshotReady reports whether a CloudStack snapshot can be used to restore a +// volume. Snapshots still being created or backed up are not ready; the +// snapshotter then calls CreateSnapshot again until they are. An empty state +// (not reported) is treated as ready to keep the previous behavior. +func isSnapshotReady(state string) bool { + switch state { + case "Allocated", "Creating", "CreatedOnPrimary", "BackingUp", "Copying": + return false + default: + return true + } +} + +// createVolumeFromSnapshot restores snapshotID into a new volume of at least sizeInGB. +func (cs *controllerServer) createVolumeFromSnapshot(ctx context.Context, req *csi.CreateVolumeRequest, name, snapshotID string, sizeInGB int64) (*csi.CreateVolumeResponse, error) { + logger := klog.FromContext(ctx) + logger.Info("Creating volume from snapshot", "snapshotID", snapshotID) + printVolumeAsJSON(req) + + snapshot, err := cs.connector.GetSnapshotByID(ctx, snapshotID) + if errors.Is(err, cloud.ErrNotFound) { + return nil, status.Errorf(codes.NotFound, "Snapshot %v not found", snapshotID) + } else if err != nil { + // Error with CloudStack + return nil, status.Errorf(codes.Internal, "Error %v", err) + } + + logger.Info("PVC created with", "size", sizeInGB) + snapshotSizeGiB := util.RoundUpBytesToGB(snapshot.Size) + if snapshotSizeGiB > sizeInGB { + logger.Info("Snapshot size is greater than the request PVC, creating volume from snapshot of size", "snapshot size:", snapshotSizeGiB) + sizeInGB = snapshotSizeGiB + } + + // A restore larger than the snapshot needs a resize afterwards. Check the + // size against CloudStack's maximum first: otherwise the snapshot would be + // copied from secondary storage on every retry, only for the resize to fail. + if sizeInGB > snapshotSizeGiB { + if maxErr := cs.checkMaxCustomDiskSize(ctx, sizeInGB); maxErr != nil { + return nil, maxErr + } + } + + volFromSnapshot, err := cs.connector.CreateVolumeFromSnapshot(ctx, snapshot.ZoneID, name, snapshot.ProjectID, snapshotID, sizeInGB) + if err != nil { + return nil, status.Errorf(codes.Internal, "Cannot create volume from snapshot %s: %v", snapshotID, err.Error()) + } + + // CloudStack creates the volume at the snapshot's size and ignores a larger + // requested size, so grow it here to satisfy the PVC request. + if growErr := cs.growVolumeFromSnapshot(ctx, volFromSnapshot, snapshotID, sizeInGB); growErr != nil { + return nil, growErr + } + + return restoredVolumeResponse(req, volFromSnapshot), nil +} + +// adoptRestoredVolume handles a CreateVolume retry for a restore whose volume +// already exists from an earlier attempt, for example after the provisioner +// timed out while CloudStack was still copying the snapshot. The volume keeps +// the snapshot's disk offering, so only size and topology are checked, and the +// content source is returned so the provisioner accepts the volume. +func (cs *controllerServer) adoptRestoredVolume(ctx context.Context, req *csi.CreateVolumeRequest, vol *cloud.Volume, snapshotID string) (*csi.CreateVolumeResponse, error) { + logger := klog.FromContext(ctx) + logger.Info("Volume from snapshot already exists, reusing it", "volumeID", vol.ID, "snapshotID", snapshotID) + + sizeInGB, err := determineSize(req) + if err != nil { + return nil, status.Error(codes.InvalidArgument, err.Error()) + } + + // An earlier attempt may have stopped before the volume was resized. + if growErr := cs.growVolumeFromSnapshot(ctx, vol, snapshotID, sizeInGB); growErr != nil { + return nil, growErr + } + + if ok, message := checkVolumeSuitable(vol, vol.DiskOfferingID, req.GetCapacityRange(), req.GetAccessibilityRequirements()); !ok { + return nil, status.Errorf(codes.AlreadyExists, "Volume %v already exists but does not satisfy request: %s", req.GetName(), message) + } + + return restoredVolumeResponse(req, vol), nil +} + +// checkMaxCustomDiskSize returns OutOfRange if sizeInGB is above CloudStack's +// maximum size for custom disk offerings. If the limit can't be read, the +// request is allowed and CloudStack decides. +func (cs *controllerServer) checkMaxCustomDiskSize(ctx context.Context, sizeInGB int64) error { + maxGB, err := cs.connector.GetMaxCustomDiskSizeGB(ctx) + if err != nil { + klog.FromContext(ctx).Error(err, "Cannot read CloudStack's maximum custom disk size; continuing") + + return nil + } + if maxGB > 0 && sizeInGB > maxGB { + return status.Errorf(codes.OutOfRange, "Requested size %d GB is above CloudStack's maximum custom disk size of %d GB", sizeInGB, maxGB) + } + + return nil +} + +func restoredVolumeResponse(req *csi.CreateVolumeRequest, vol *cloud.Volume) *csi.CreateVolumeResponse { + return &csi.CreateVolumeResponse{ + Volume: &csi.Volume{ + VolumeId: vol.ID, + CapacityBytes: vol.Size, + VolumeContext: req.GetParameters(), + ContentSource: req.GetVolumeContentSource(), + AccessibleTopology: []*csi.Topology{ + Topology{ZoneID: vol.ZoneID}.ToCSI(), + }, + }, + } +} + +// growVolumeFromSnapshot resizes a volume created from a snapshot when it is +// smaller than the requested size. On failure the undersized volume is deleted +// so the next CreateVolume retry starts clean. +func (cs *controllerServer) growVolumeFromSnapshot(ctx context.Context, vol *cloud.Volume, snapshotID string, sizeInGB int64) error { + logger := klog.FromContext(ctx) + + requiredBytes := util.GigaBytesToBytes(sizeInGB) + if vol.Size >= requiredBytes { + return nil + } + + logger.Info("Volume from snapshot is smaller than requested, resizing", + "volumeID", vol.ID, "currentBytes", vol.Size, "requestedGB", sizeInGB) + + if expandErr := cs.connector.ExpandVolume(ctx, vol.ID, sizeInGB); expandErr != nil { + if delErr := cs.connector.DeleteVolume(ctx, vol.ID); delErr != nil { + logger.Error(delErr, "Failed to delete undersized volume after resize failure", "volumeID", vol.ID) + } + + return status.Errorf(codes.Internal, "Cannot resize volume %s created from snapshot %s to %d GB: %v", + vol.ID, snapshotID, sizeInGB, expandErr) + } + vol.Size = requiredBytes + + return nil +} + func determineSize(req *csi.CreateVolumeRequest) (int64, error) { var sizeInGB int64 @@ -341,11 +447,41 @@ func (cs *controllerServer) DeleteVolume(ctx context.Context, req *csi.DeleteVol ) err := cs.connector.DeleteVolume(ctx, volumeID) - if err != nil && !errors.Is(err, cloud.ErrNotFound) { - return nil, status.Errorf(codes.Internal, "Cannot delete volume %s: %s", volumeID, err.Error()) + if err == nil || errors.Is(err, cloud.ErrNotFound) { + return &csi.DeleteVolumeResponse{}, nil + } + + // CloudStack uses the same error code for "no such volume" and for refusals + // such as "volume is attached". Only treat the delete as done if the volume + // is really gone; otherwise report why, so nothing is left behind unnoticed. + vol, getErr := cs.connector.GetVolumeByID(ctx, volumeID) + switch { + case errors.Is(getErr, cloud.ErrNotFound): + return &csi.DeleteVolumeResponse{}, nil + case getErr != nil: + return nil, status.Errorf(codes.Internal, "Cannot delete volume %s: %v (checking whether it still exists failed: %v)", volumeID, err, getErr) + case isVolumeDestroyed(vol.State): + logger.Info("Volume already destroyed in CloudStack; it will be expunged by CloudStack", + "volumeID", volumeID, "state", vol.State) + + return &csi.DeleteVolumeResponse{}, nil + case vol.VirtualMachineID != "": + return nil, status.Errorf(codes.FailedPrecondition, "Cannot delete volume %s: it is attached to VM %s: %v", + volumeID, vol.VirtualMachineID, err) + default: + return nil, status.Errorf(codes.Internal, "Cannot delete volume %s: %v", volumeID, err) } +} - return &csi.DeleteVolumeResponse{}, nil +// isVolumeDestroyed reports whether a CloudStack volume has been destroyed and +// only waits for CloudStack to expunge it. +func isVolumeDestroyed(state string) bool { + switch state { + case "Destroy", "Expunging", "Expunged": + return true + default: + return false + } } func (cs *controllerServer) CreateSnapshot(ctx context.Context, req *csi.CreateSnapshotRequest) (*csi.CreateSnapshotResponse, error) { @@ -373,11 +509,37 @@ func (cs *controllerServer) CreateSnapshot(ctx context.Context, req *csi.CreateS } klog.V(4).Infof("CreateSnapshot of volume: %s", volume.ID) - snapshot, err := cs.connector.CreateSnapshot(ctx, volume.ID, req.GetName()) - if errors.Is(err, cloud.ErrAlreadyExists) { - return nil, status.Errorf(codes.AlreadyExists, "Snapshot name conflict: already exists for a different source volume") - } else if err != nil { - return nil, status.Errorf(codes.Internal, "Failed to create snapshot for volume %s: %v", volume.ID, err.Error()) + + // Serialize calls for the same snapshot name. A timed-out call keeps running + // against CloudStack, so a retry must not race it and create a duplicate. + if acquired := cs.volumeLocks.TryAcquire(req.GetName()); !acquired { + return nil, status.Errorf(codes.Aborted, "An operation for snapshot %s is already in progress", req.GetName()) + } + defer cs.volumeLocks.Release(req.GetName()) + + // CreateSnapshot must be idempotent (CSI spec): the snapshotter retries with + // the same name after timeouts, so return an existing snapshot instead of + // creating a duplicate in CloudStack. + snapshot, err := cs.connector.GetSnapshotByName(ctx, req.GetName()) + switch { + case err == nil: + if snapshot.VolumeID != volume.ID { + return nil, status.Errorf(codes.AlreadyExists, "Snapshot %q already exists for a different source volume", req.GetName()) + } + klog.V(4).Infof("CreateSnapshot: snapshot %s already exists for volume %s, returning it", snapshot.ID, volume.ID) + case errors.Is(err, cloud.ErrNotFound): + snapshot, err = cs.connector.CreateSnapshot(ctx, volume.ID, req.GetName()) + if errors.Is(err, cloud.ErrAlreadyExists) { + return nil, status.Errorf(codes.AlreadyExists, "Snapshot name conflict: already exists for a different source volume") + } else if err != nil { + return nil, status.Errorf(codes.Internal, "Failed to create snapshot for volume %s: %v", volume.ID, err.Error()) + } + default: + return nil, status.Errorf(codes.Internal, "Failed to look up snapshot %q: %v", req.GetName(), err) + } + + if snapshot.State == "Error" { + return nil, status.Errorf(codes.Internal, "Snapshot %s of volume %s is in Error state in CloudStack", snapshot.ID, volume.ID) } t, err := time.Parse("2006-01-02T15:04:05-0700", snapshot.CreatedAt) @@ -387,12 +549,20 @@ func (cs *controllerServer) CreateSnapshot(ctx context.Context, req *csi.CreateS ts := timestamppb.New(t) + // Report the restore size so Kubernetes can show and enforce it. + // Fall back to the source volume size if CloudStack returns no virtualsize. + sizeBytes := snapshot.Size + if sizeBytes == 0 { + sizeBytes = volume.Size + } + resp := &csi.CreateSnapshotResponse{ Snapshot: &csi.Snapshot{ SnapshotId: snapshot.ID, SourceVolumeId: volume.ID, + SizeBytes: sizeBytes, CreationTime: ts, - ReadyToUse: true, + ReadyToUse: isSnapshotReady(snapshot.State), }, } @@ -434,8 +604,9 @@ func (cs *controllerServer) ListSnapshots(ctx context.Context, req *csi.ListSnap Snapshot: &csi.Snapshot{ SnapshotId: snap.ID, SourceVolumeId: snap.VolumeID, + SizeBytes: snap.Size, CreationTime: ts, - ReadyToUse: true, + ReadyToUse: isSnapshotReady(snap.State), }, } entries = append(entries, entry) @@ -454,14 +625,24 @@ func (cs *controllerServer) DeleteSnapshot(ctx context.Context, req *csi.DeleteS klog.V(4).Infof("DeleteSnapshot for snapshotID: %s", snapshotID) err := cs.connector.DeleteSnapshot(ctx, snapshotID) - if errors.Is(err, cloud.ErrNotFound) { + if err == nil || errors.Is(err, cloud.ErrNotFound) { // Per CSI spec, return OK if snapshot does not exist return &csi.DeleteSnapshotResponse{}, nil - } else if err != nil { - return nil, status.Errorf(codes.Internal, "Error %v", err) } - return &csi.DeleteSnapshotResponse{}, nil + // CloudStack reports an already deleted snapshot with several different + // messages. Treat the delete as done only if the snapshot is really gone. + snap, getErr := cs.connector.GetSnapshotByID(ctx, snapshotID) + switch { + case errors.Is(getErr, cloud.ErrNotFound): + return &csi.DeleteSnapshotResponse{}, nil + case getErr != nil: + return nil, status.Errorf(codes.Internal, "Cannot delete snapshot %s: %v (checking whether it still exists failed: %v)", snapshotID, err, getErr) + case snap.State == "Destroyed": + return &csi.DeleteSnapshotResponse{}, nil + default: + return nil, status.Errorf(codes.Internal, "Cannot delete snapshot %s (state %s): %v", snapshotID, snap.State, err) + } } func (cs *controllerServer) ControllerPublishVolume(ctx context.Context, req *csi.ControllerPublishVolumeRequest) (*csi.ControllerPublishVolumeResponse, error) { diff --git a/pkg/driver/controller_delete_test.go b/pkg/driver/controller_delete_test.go new file mode 100644 index 0000000..182b063 --- /dev/null +++ b/pkg/driver/controller_delete_test.go @@ -0,0 +1,116 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package driver + +import ( + "context" + "errors" + "testing" + + "github.com/container-storage-interface/spec/lib/go/csi" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud" + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud/fake" +) + +// deleteTestConnector makes DeleteVolume/DeleteSnapshot fail with a given error +// and controls what the follow-up lookup finds (nil means "not found"). +type deleteTestConnector struct { + cloud.Interface + + deleteErr error + volume *cloud.Volume + snapshot *cloud.Snapshot +} + +func (c *deleteTestConnector) DeleteVolume(context.Context, string) error { return c.deleteErr } + +func (c *deleteTestConnector) GetVolumeByID(context.Context, string) (*cloud.Volume, error) { + if c.volume == nil { + return nil, cloud.ErrNotFound + } + + return c.volume, nil +} + +func (c *deleteTestConnector) DeleteSnapshot(context.Context, string) error { return c.deleteErr } + +func (c *deleteTestConnector) GetSnapshotByID(context.Context, string) (*cloud.Snapshot, error) { + if c.snapshot == nil { + return nil, cloud.ErrNotFound + } + + return c.snapshot, nil +} + +func TestDeleteVolumeVerifiesFailedDeletes(t *testing.T) { + refused := errors.New("CloudStack API error 431 (CSExceptionErrorCode: 4350): Please specify a volume that is not attached to any VM") + + cases := []struct { + name string + deleteErr error + volume *cloud.Volume + want codes.Code + }{ + {"deleted", nil, nil, codes.OK}, + {"already not found", cloud.ErrNotFound, nil, codes.OK}, + {"error but volume gone", refused, nil, codes.OK}, + {"error and volume destroyed", refused, &cloud.Volume{ID: "v", State: "Destroy"}, codes.OK}, + {"error and volume attached", refused, &cloud.Volume{ID: "v", State: "Ready", VirtualMachineID: "vm-1"}, codes.FailedPrecondition}, + {"error and volume still there", refused, &cloud.Volume{ID: "v", State: "Ready"}, codes.Internal}, + } + for _, testCase := range cases { + t.Run(testCase.name, func(t *testing.T) { + cs := NewControllerServer(&deleteTestConnector{Interface: fake.New(), deleteErr: testCase.deleteErr, volume: testCase.volume}) + _, err := cs.DeleteVolume(context.Background(), &csi.DeleteVolumeRequest{VolumeId: "v"}) + if got := status.Code(err); got != testCase.want { + t.Errorf("DeleteVolume() code = %v, want %v (err: %v)", got, testCase.want, err) + } + }) + } +} + +func TestDeleteSnapshotVerifiesFailedDeletes(t *testing.T) { + cases := []struct { + name string + deleteErr error + snapshot *cloud.Snapshot + want codes.Code + }{ + {"deleted", nil, nil, codes.OK}, + {"already not found", cloud.ErrNotFound, nil, codes.OK}, + {"unknown id", errors.New("CloudStack API error 431 (CSExceptionErrorCode: 4350): Invalid parameter id"), nil, codes.OK}, + {"already destroyed", errors.New(`{"errorcode":431,"errortext":"Snapshot [...] is already destroyed"}`), nil, codes.OK}, + {"removed record", errors.New(`{"errorcode":431,"errortext":"unable to find a snapshot with id 33"}`), nil, codes.OK}, + {"error and snapshot destroyed", errors.New("failed"), &cloud.Snapshot{ID: "s", State: "Destroyed"}, codes.OK}, + {"error and snapshot still there", errors.New("failed"), &cloud.Snapshot{ID: "s", State: "BackedUp"}, codes.Internal}, + } + for _, testCase := range cases { + t.Run(testCase.name, func(t *testing.T) { + cs := NewControllerServer(&deleteTestConnector{Interface: fake.New(), deleteErr: testCase.deleteErr, snapshot: testCase.snapshot}) + _, err := cs.DeleteSnapshot(context.Background(), &csi.DeleteSnapshotRequest{SnapshotId: "s"}) + if got := status.Code(err); got != testCase.want { + t.Errorf("DeleteSnapshot() code = %v, want %v (err: %v)", got, testCase.want, err) + } + }) + } +} diff --git a/pkg/driver/controller_restore_resize_test.go b/pkg/driver/controller_restore_resize_test.go new file mode 100644 index 0000000..688edc6 --- /dev/null +++ b/pkg/driver/controller_restore_resize_test.go @@ -0,0 +1,57 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package driver + +import ( + "context" + "testing" + + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud" + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud/fake" + "github.com/cloudstack/cloudstack-csi-driver/pkg/util" +) + +func TestGrowVolumeFromSnapshot(t *testing.T) { + cs, ok := NewControllerServer(fake.New()).(*controllerServer) + if !ok { + t.Fatal("unexpected controller server type") + } + + // Volume known to the fake connector, smaller than requested. + vol := &cloud.Volume{ID: "ace9f28b-3081-40c1-8353-4cc3e3014072", Size: 10} + if err := cs.growVolumeFromSnapshot(context.Background(), vol, "snap-1", 5); err != nil { + t.Fatalf("growVolumeFromSnapshot failed: %v", err) + } + if want := util.GigaBytesToBytes(5); vol.Size != want { + t.Errorf("expected size %d, got %d", want, vol.Size) + } + + // Already large enough: no-op. + big := &cloud.Volume{ID: "does-not-matter", Size: util.GigaBytesToBytes(10)} + if err := cs.growVolumeFromSnapshot(context.Background(), big, "snap-1", 5); err != nil { + t.Errorf("expected no-op for large enough volume, got %v", err) + } + + // Resize failure (unknown volume) returns an error. + missing := &cloud.Volume{ID: "unknown", Size: 10} + if err := cs.growVolumeFromSnapshot(context.Background(), missing, "snap-1", 5); err == nil { + t.Error("expected error when resize fails") + } +} diff --git a/pkg/driver/controller_restore_retry_test.go b/pkg/driver/controller_restore_retry_test.go new file mode 100644 index 0000000..f2d0ac3 --- /dev/null +++ b/pkg/driver/controller_restore_retry_test.go @@ -0,0 +1,135 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package driver + +import ( + "context" + "testing" + + "github.com/container-storage-interface/spec/lib/go/csi" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud" + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud/fake" + "github.com/cloudstack/cloudstack-csi-driver/pkg/util" +) + +const restoreTestSourceVolume = "ace9f28b-3081-40c1-8353-4cc3e3014072" + +// maxSizeConnector reports a maximum custom disk size and counts restores. +type maxSizeConnector struct { + cloud.Interface + + maxGB int64 + restores int +} + +func (c *maxSizeConnector) GetMaxCustomDiskSizeGB(context.Context) (int64, error) { + return c.maxGB, nil +} + +func (c *maxSizeConnector) CreateVolumeFromSnapshot(ctx context.Context, zoneID, name, projectID, snapshotID string, sizeInGB int64) (*cloud.Volume, error) { + c.restores++ + + return c.Interface.CreateVolumeFromSnapshot(ctx, zoneID, name, projectID, snapshotID, sizeInGB) +} + +func restoreRequest(t *testing.T, cs csi.ControllerServer, name string, sizeGB int64) *csi.CreateVolumeRequest { + t.Helper() + + snap, err := cs.CreateSnapshot(context.Background(), &csi.CreateSnapshotRequest{ + Name: "snap-for-" + name, + SourceVolumeId: restoreTestSourceVolume, + }) + if err != nil { + t.Fatalf("CreateSnapshot failed: %v", err) + } + + return &csi.CreateVolumeRequest{ + Name: name, + CapacityRange: &csi.CapacityRange{RequiredBytes: util.GigaBytesToBytes(sizeGB)}, + VolumeCapabilities: []*csi.VolumeCapability{{ + AccessType: &csi.VolumeCapability_Mount{Mount: &csi.VolumeCapability_MountVolume{}}, + AccessMode: &csi.VolumeCapability_AccessMode{Mode: csi.VolumeCapability_AccessMode_SINGLE_NODE_WRITER}, + }}, + // Differs from the disk offering the restored volume gets. + Parameters: map[string]string{DiskOfferingKey: "storageclass-offering"}, + VolumeContentSource: &csi.VolumeContentSource{ + Type: &csi.VolumeContentSource_Snapshot{ + Snapshot: &csi.VolumeContentSource_SnapshotSource{SnapshotId: snap.GetSnapshot().GetSnapshotId()}, + }, + }, + } +} + +func TestCreateVolumeRetryReusesRestoredVolume(t *testing.T) { + cs := NewControllerServer(fake.New()) + req := restoreRequest(t, cs, "pvc-retry", 1) + + first, err := cs.CreateVolume(context.Background(), req) + if err != nil { + t.Fatalf("first CreateVolume failed: %v", err) + } + + // Retry with the same name, as the provisioner does after a timeout. + second, err := cs.CreateVolume(context.Background(), req) + if err != nil { + t.Fatalf("retried CreateVolume failed: %v", err) + } + if second.GetVolume().GetVolumeId() != first.GetVolume().GetVolumeId() { + t.Errorf("retry returned volume %q, want the existing %q", second.GetVolume().GetVolumeId(), first.GetVolume().GetVolumeId()) + } + if second.GetVolume().GetContentSource() == nil { + t.Error("retry did not return the content source; the provisioner would delete the volume") + } +} + +func TestCreateVolumeRetryGrowsUndersizedRestoredVolume(t *testing.T) { + cs := NewControllerServer(fake.New()) + req := restoreRequest(t, cs, "pvc-grow", 1) + if _, err := cs.CreateVolume(context.Background(), req); err != nil { + t.Fatalf("first CreateVolume failed: %v", err) + } + + // The retry asks for more than the earlier attempt created. + req.CapacityRange = &csi.CapacityRange{RequiredBytes: util.GigaBytesToBytes(3)} + resp, err := cs.CreateVolume(context.Background(), req) + if err != nil { + t.Fatalf("retried CreateVolume failed: %v", err) + } + if got, want := resp.GetVolume().GetCapacityBytes(), util.GigaBytesToBytes(3); got != want { + t.Errorf("capacity = %d, want %d", got, want) + } +} + +func TestCreateVolumeFromSnapshotAboveMaxSizeFailsBeforeRestoring(t *testing.T) { + conn := &maxSizeConnector{Interface: fake.New(), maxGB: 1024} + cs := NewControllerServer(conn) + req := restoreRequest(t, cs, "pvc-huge", 2000) + + _, err := cs.CreateVolume(context.Background(), req) + if got := status.Code(err); got != codes.OutOfRange { + t.Errorf("CreateVolume code = %v, want OutOfRange (err: %v)", got, err) + } + if conn.restores != 0 { + t.Errorf("snapshot was restored %d times; want 0 when the size is above the maximum", conn.restores) + } +} diff --git a/pkg/driver/controller_snapshot_idempotent_test.go b/pkg/driver/controller_snapshot_idempotent_test.go new file mode 100644 index 0000000..f196276 --- /dev/null +++ b/pkg/driver/controller_snapshot_idempotent_test.go @@ -0,0 +1,50 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package driver + +import ( + "context" + "testing" + + "github.com/container-storage-interface/spec/lib/go/csi" + + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud/fake" +) + +func TestCreateSnapshotIsIdempotent(t *testing.T) { + cs := NewControllerServer(fake.New()) + req := &csi.CreateSnapshotRequest{ + Name: "snap-idem", + SourceVolumeId: "ace9f28b-3081-40c1-8353-4cc3e3014072", + } + + first, err := cs.CreateSnapshot(context.Background(), req) + if err != nil { + t.Fatalf("first CreateSnapshot failed: %v", err) + } + second, err := cs.CreateSnapshot(context.Background(), req) + if err != nil { + t.Fatalf("second CreateSnapshot failed: %v", err) + } + if first.GetSnapshot().GetSnapshotId() != second.GetSnapshot().GetSnapshotId() { + t.Errorf("expected same snapshot ID on retry, got %q and %q", + first.GetSnapshot().GetSnapshotId(), second.GetSnapshot().GetSnapshotId()) + } +} diff --git a/pkg/driver/controller_snapshot_ready_test.go b/pkg/driver/controller_snapshot_ready_test.go new file mode 100644 index 0000000..fd7622a --- /dev/null +++ b/pkg/driver/controller_snapshot_ready_test.go @@ -0,0 +1,39 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package driver + +import "testing" + +func TestIsSnapshotReady(t *testing.T) { + cases := map[string]bool{ + "Allocated": false, + "Creating": false, + "CreatedOnPrimary": false, + "BackingUp": false, + "Copying": false, + "BackedUp": true, + "": true, + } + for state, want := range cases { + if got := isSnapshotReady(state); got != want { + t.Errorf("isSnapshotReady(%q) = %v, want %v", state, got, want) + } + } +} diff --git a/pkg/driver/controller_snapshot_size_test.go b/pkg/driver/controller_snapshot_size_test.go new file mode 100644 index 0000000..a1a1715 --- /dev/null +++ b/pkg/driver/controller_snapshot_size_test.go @@ -0,0 +1,43 @@ +// +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. +// + +package driver + +import ( + "context" + "testing" + + "github.com/container-storage-interface/spec/lib/go/csi" + + "github.com/cloudstack/cloudstack-csi-driver/pkg/cloud/fake" +) + +func TestCreateSnapshotReportsSize(t *testing.T) { + cs := NewControllerServer(fake.New()) + resp, err := cs.CreateSnapshot(context.Background(), &csi.CreateSnapshotRequest{ + Name: "snap-size", + SourceVolumeId: "ace9f28b-3081-40c1-8353-4cc3e3014072", + }) + if err != nil { + t.Fatalf("CreateSnapshot failed: %v", err) + } + if resp.GetSnapshot().GetSizeBytes() == 0 { + t.Errorf("expected non-zero SizeBytes, got 0") + } +}