From da874edd1f44e6e1ec3e1ae76fd1d44fb7f25033 Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Fri, 2 Oct 2026 18:16:23 +0200 Subject: [PATCH 1/8] Report snapshot size in CreateSnapshot and ListSnapshots CreateSnapshotResponse and ListSnapshots entries never set SizeBytes, so VolumeSnapshot.status.restoreSize was always 0. Kubernetes uses this value to show the minimum restore size and to reject PVCs that are too small. Use the snapshot virtualsize from CloudStack, falling back to the source volume size. Also populate Size/CreatedAt in GetSnapshotByID and Size in GetSnapshotByName so callers see the real snapshot size. --- pkg/cloud/snapshots.go | 3 ++ pkg/driver/controller.go | 9 +++++ pkg/driver/controller_snapshot_size_test.go | 43 +++++++++++++++++++++ 3 files changed, 55 insertions(+) create mode 100644 pkg/driver/controller_snapshot_size_test.go diff --git a/pkg/cloud/snapshots.go b/pkg/cloud/snapshots.go index fae632a..8c77e2d 100644 --- a/pkg/cloud/snapshots.go +++ b/pkg/cloud/snapshots.go @@ -46,10 +46,12 @@ func (c *client) GetSnapshotByID(ctx context.Context, snapshotID string) (*Snaps return &Snapshot{ ID: snapshot.Id, Name: snapshot.Name, + Size: snapshot.Virtualsize, DomainID: snapshot.Domainid, ProjectID: snapshot.Projectid, ZoneID: snapshot.Zoneid, VolumeID: snapshot.Volumeid, + CreatedAt: snapshot.Created, }, nil } @@ -107,6 +109,7 @@ func (c *client) GetSnapshotByName(ctx context.Context, name string) (*Snapshot, return &Snapshot{ ID: snapshot.Id, Name: snapshot.Name, + Size: snapshot.Virtualsize, DomainID: snapshot.Domainid, ProjectID: snapshot.Projectid, ZoneID: snapshot.Zoneid, diff --git a/pkg/driver/controller.go b/pkg/driver/controller.go index 9ca1185..96f5164 100644 --- a/pkg/driver/controller.go +++ b/pkg/driver/controller.go @@ -387,10 +387,18 @@ 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, }, @@ -434,6 +442,7 @@ 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, }, 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") + } +} From e9049bb877be156cd9769e92b0a12e778b3c9816 Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Fri, 2 Oct 2026 18:16:23 +0200 Subject: [PATCH 2/8] Resize volumes restored from snapshot to the requested size CloudStack creates a volume from a snapshot at the snapshot's size and ignores a larger requested size. The driver returned that smaller capacity, so external-provisioner rejected it with 'created volume capacity X less than requested capacity Y', deleted the volume and retried forever. After creating the volume, resize it to the requested size when it is smaller. If the resize fails, delete the undersized volume and return an error. The node plugin already grows the filesystem on NodeStage. --- pkg/driver/controller.go | 33 ++++++++++++ pkg/driver/controller_restore_resize_test.go | 57 ++++++++++++++++++++ 2 files changed, 90 insertions(+) create mode 100644 pkg/driver/controller_restore_resize_test.go diff --git a/pkg/driver/controller.go b/pkg/driver/controller.go index 96f5164..f9f0d94 100644 --- a/pkg/driver/controller.go +++ b/pkg/driver/controller.go @@ -174,6 +174,12 @@ func (cs *controllerServer) CreateVolume(ctx context.Context, req *csi.CreateVol 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 + } + resp := &csi.CreateVolumeResponse{ Volume: &csi.Volume{ VolumeId: volFromSnapshot.ID, @@ -285,6 +291,33 @@ func checkVolumeSuitable(vol *cloud.Volume, return true, "" } +// 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 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") + } +} From d4c097d7311d834d7ef21315369050dd5e3f89fb Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Fri, 2 Oct 2026 18:16:23 +0200 Subject: [PATCH 3/8] Make CreateSnapshot idempotent CreateSnapshot always called CloudStack createSnapshot, so a retry from csi-snapshotter (e.g. after a timeout) created a duplicate snapshot. The CSI spec requires CreateSnapshot to return the existing snapshot when one with the same name and source volume already exists. Look up the snapshot by name first; return it if it belongs to the same volume, return AlreadyExists if it belongs to a different volume, and only create a new snapshot when none exists. --- pkg/driver/controller.go | 24 +++++++-- .../controller_snapshot_idempotent_test.go | 50 +++++++++++++++++++ 2 files changed, 69 insertions(+), 5 deletions(-) create mode 100644 pkg/driver/controller_snapshot_idempotent_test.go diff --git a/pkg/driver/controller.go b/pkg/driver/controller.go index f9f0d94..cd60a34 100644 --- a/pkg/driver/controller.go +++ b/pkg/driver/controller.go @@ -406,11 +406,25 @@ 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()) + // 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) } t, err := time.Parse("2006-01-02T15:04:05-0700", snapshot.CreatedAt) 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()) + } +} From a016378737807d6252a791e0530174007de02be6 Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Fri, 2 Oct 2026 18:16:23 +0200 Subject: [PATCH 4/8] Move snapshot restore out of CreateVolume --- pkg/driver/controller.go | 91 ++++++++++++++++++++-------------------- 1 file changed, 46 insertions(+), 45 deletions(-) diff --git a/pkg/driver/controller.go b/pkg/driver/controller.go index cd60a34..eed7584 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) @@ -148,51 +147,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()) - } - - // 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 - } - - 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. @@ -291,6 +247,51 @@ func checkVolumeSuitable(vol *cloud.Volume, 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 + } + + 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 &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(), + }, + }, + }, nil +} + // 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. From bef97b52770fdfc6d061640f30a3ad48bbb47421 Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Fri, 2 Oct 2026 18:16:23 +0200 Subject: [PATCH 5/8] Serialise CreateSnapshot per name and report ReadyToUse from snapshot state A timed-out CreateSnapshot keeps running against CloudStack, so a retry could miss the in-progress snapshot and create a duplicate. Hold a per-name lock for the call and return Aborted to concurrent calls. CreateSnapshot and ListSnapshots always returned ReadyToUse: true, even while CloudStack was still creating or backing up the snapshot. Report ReadyToUse=false until the snapshot leaves Allocated/Creating/ CreatedOnPrimary/BackingUp/Copying, and an error for the Error state. --- pkg/cloud/cloud.go | 3 ++ pkg/cloud/snapshots.go | 4 ++ pkg/driver/controller.go | 29 ++++++++++++++- pkg/driver/controller_snapshot_ready_test.go | 39 ++++++++++++++++++++ 4 files changed, 73 insertions(+), 2 deletions(-) create mode 100644 pkg/driver/controller_snapshot_ready_test.go diff --git a/pkg/cloud/cloud.go b/pkg/cloud/cloud.go index bed6aa8..456b370 100644 --- a/pkg/cloud/cloud.go +++ b/pkg/cloud/cloud.go @@ -83,6 +83,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/snapshots.go b/pkg/cloud/snapshots.go index 8c77e2d..f6deb03 100644 --- a/pkg/cloud/snapshots.go +++ b/pkg/cloud/snapshots.go @@ -46,6 +46,7 @@ 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, @@ -72,6 +73,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, @@ -109,6 +111,7 @@ 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, @@ -151,6 +154,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/driver/controller.go b/pkg/driver/controller.go index eed7584..7304572 100644 --- a/pkg/driver/controller.go +++ b/pkg/driver/controller.go @@ -247,6 +247,19 @@ 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) @@ -407,6 +420,14 @@ func (cs *controllerServer) CreateSnapshot(ctx context.Context, req *csi.CreateS } klog.V(4).Infof("CreateSnapshot of volume: %s", volume.ID) + + // 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. @@ -428,6 +449,10 @@ func (cs *controllerServer) CreateSnapshot(ctx context.Context, req *csi.CreateS 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) if err != nil { return nil, status.Errorf(codes.Internal, "Failed to parse snapshot creation time: %v", err) @@ -448,7 +473,7 @@ func (cs *controllerServer) CreateSnapshot(ctx context.Context, req *csi.CreateS SourceVolumeId: volume.ID, SizeBytes: sizeBytes, CreationTime: ts, - ReadyToUse: true, + ReadyToUse: isSnapshotReady(snapshot.State), }, } @@ -492,7 +517,7 @@ func (cs *controllerServer) ListSnapshots(ctx context.Context, req *csi.ListSnap SourceVolumeId: snap.VolumeID, SizeBytes: snap.Size, CreationTime: ts, - ReadyToUse: true, + ReadyToUse: isSnapshotReady(snap.State), }, } entries = append(entries, entry) 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) + } + } +} From ee5d32cc23a4b103797f66ccf4caaa6953ca261d Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Fri, 2 Oct 2026 18:16:23 +0200 Subject: [PATCH 6/8] Treat already-destroyed snapshots as deleted in DeleteSnapshot If a snapshot was destroyed outside the driver, or an earlier delete completed after the RPC timed out, CloudStack returns errorcode 431 '... is already destroyed'. That was not mapped to ErrNotFound, so DeleteSnapshot returned Internal and the VolumeSnapshotContent (and its VolumeSnapshot) could never be deleted. The CSI spec requires DeleteSnapshot to succeed when the snapshot no longer exists. --- pkg/cloud/snapshots.go | 16 ++++++++++++-- pkg/cloud/snapshots_test.go | 42 +++++++++++++++++++++++++++++++++++++ 2 files changed, 56 insertions(+), 2 deletions(-) create mode 100644 pkg/cloud/snapshots_test.go diff --git a/pkg/cloud/snapshots.go b/pkg/cloud/snapshots.go index f6deb03..74525f0 100644 --- a/pkg/cloud/snapshots.go +++ b/pkg/cloud/snapshots.go @@ -86,14 +86,26 @@ 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 + if isSnapshotGoneError(err) { return ErrNotFound } return err } +// isSnapshotGoneError reports whether a CloudStack deleteSnapshot error means +// the snapshot no longer exists, so DeleteSnapshot can succeed idempotently. +func isSnapshotGoneError(err error) bool { + if err == nil { + return false + } + msg := err.Error() + + // 4350: InvalidParameterValueException, e.g. unknown snapshot ID. + // "already destroyed": snapshot was deleted earlier (errorcode 431). + return strings.Contains(msg, "4350") || strings.Contains(msg, "is already destroyed") +} + func (c *client) GetSnapshotByName(ctx context.Context, name string) (*Snapshot, error) { logger := klog.FromContext(ctx) logger.V(2).Info("CloudStack API call", "command", "GetSnapshotByName", "params", map[string]string{ diff --git a/pkg/cloud/snapshots_test.go b/pkg/cloud/snapshots_test.go new file mode 100644 index 0000000..f77a559 --- /dev/null +++ b/pkg/cloud/snapshots_test.go @@ -0,0 +1,42 @@ +// +// 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 ( + "errors" + "testing" +) + +func TestIsSnapshotGoneError(t *testing.T) { + cases := []struct { + err error + want bool + }{ + {nil, false}, + {errors.New(`CloudStack API error 431 (CSExceptionErrorCode: 4350): Invalid parameter id`), true}, + {errors.New(`Undefined error: {"errorcode":431,"errortext":"Snapshot [Snapshot {\"id\":21,\"state\":\"Destroyed\"}] is already destroyed"}`), true}, + {errors.New(`CloudStack API error 530: Failed to delete snapshot`), false}, + } + for _, c := range cases { + if got := isSnapshotGoneError(c.err); got != c.want { + t.Errorf("isSnapshotGoneError(%v) = %v, want %v", c.err, got, c.want) + } + } +} From 8725abd0ce312b4b752c6199b810989014a15a8d Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Tue, 6 Oct 2026 12:42:01 +0200 Subject: [PATCH 7/8] Verify that volumes and snapshots are gone before treating a failed delete as done CloudStack returns error 4350 (InvalidParameterValueException) both for an unknown ID and for refusals such as deleting a volume that is still attached. DeleteVolume mapped every 4350 to ErrNotFound, so a refused delete was reported as success: Kubernetes removed the PV and the CloudStack volume was left behind with nothing tracking it. DeleteSnapshot matched error text instead, which missed messages such as 'unable to find a snapshot with id N', leaving VolumeSnapshots stuck. The cloud client now returns delete errors unchanged. After a failed delete, the controller looks the volume or snapshot up by ID and only reports success if it is gone or already destroyed. A volume that is still attached returns FailedPrecondition; anything else returns the CloudStack error. --- pkg/cloud/cloud.go | 3 + pkg/cloud/snapshots.go | 20 +---- pkg/cloud/snapshots_test.go | 42 ---------- pkg/cloud/volumes.go | 10 +-- pkg/driver/controller.go | 54 +++++++++++-- pkg/driver/controller_delete_test.go | 116 +++++++++++++++++++++++++++ 6 files changed, 174 insertions(+), 71 deletions(-) delete mode 100644 pkg/cloud/snapshots_test.go create mode 100644 pkg/driver/controller_delete_test.go diff --git a/pkg/cloud/cloud.go b/pkg/cloud/cloud.go index 456b370..56551c7 100644 --- a/pkg/cloud/cloud.go +++ b/pkg/cloud/cloud.go @@ -69,6 +69,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. diff --git a/pkg/cloud/snapshots.go b/pkg/cloud/snapshots.go index 74525f0..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" @@ -86,26 +85,13 @@ 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 isSnapshotGoneError(err) { - 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 } -// isSnapshotGoneError reports whether a CloudStack deleteSnapshot error means -// the snapshot no longer exists, so DeleteSnapshot can succeed idempotently. -func isSnapshotGoneError(err error) bool { - if err == nil { - return false - } - msg := err.Error() - - // 4350: InvalidParameterValueException, e.g. unknown snapshot ID. - // "already destroyed": snapshot was deleted earlier (errorcode 431). - return strings.Contains(msg, "4350") || strings.Contains(msg, "is already destroyed") -} - func (c *client) GetSnapshotByName(ctx context.Context, name string) (*Snapshot, error) { logger := klog.FromContext(ctx) logger.V(2).Info("CloudStack API call", "command", "GetSnapshotByName", "params", map[string]string{ diff --git a/pkg/cloud/snapshots_test.go b/pkg/cloud/snapshots_test.go deleted file mode 100644 index f77a559..0000000 --- a/pkg/cloud/snapshots_test.go +++ /dev/null @@ -1,42 +0,0 @@ -// -// 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 ( - "errors" - "testing" -) - -func TestIsSnapshotGoneError(t *testing.T) { - cases := []struct { - err error - want bool - }{ - {nil, false}, - {errors.New(`CloudStack API error 431 (CSExceptionErrorCode: 4350): Invalid parameter id`), true}, - {errors.New(`Undefined error: {"errorcode":431,"errortext":"Snapshot [Snapshot {\"id\":21,\"state\":\"Destroyed\"}] is already destroyed"}`), true}, - {errors.New(`CloudStack API error 530: Failed to delete snapshot`), false}, - } - for _, c := range cases { - if got := isSnapshotGoneError(c.err); got != c.want { - t.Errorf("isSnapshotGoneError(%v) = %v, want %v", c.err, got, c.want) - } - } -} 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 7304572..a554a84 100644 --- a/pkg/driver/controller.go +++ b/pkg/driver/controller.go @@ -388,11 +388,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 } - 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) + } +} + +// 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) { @@ -536,14 +566,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) + } + }) + } +} From ab65e6455c558106660ca0bdc30ec7a7cc331dc4 Mon Sep 17 00:00:00 2001 From: wilsom10 <26769729+mw-0@users.noreply.github.com> Date: Tue, 6 Oct 2026 14:10:03 +0200 Subject: [PATCH 8/8] Reuse an existing restored volume on CreateVolume retries and check the maximum size before restoring When a restore from snapshot took longer than the provisioner's timeout, the retry found the volume by name and returned it without a content source and after comparing it with the StorageClass's disk offering. The restored volume keeps the snapshot's offering, and external-provisioner deletes a restored volume without a content source ('volume content source missing'), so every retry copied the snapshot from secondary storage again and the PVC never bound. A retry for a restore now reuses the existing volume: it is grown if an earlier attempt stopped before the resize, checked for size and topology only, and returned with its content source. A restore larger than the snapshot is now checked against CloudStack's maximum custom disk size (listCapabilities) before anything is created. Previously each retry restored the whole snapshot only for the resize to fail; now the request fails at once with OutOfRange. --- pkg/cloud/capabilities.go | 41 ++++++ pkg/cloud/cloud.go | 4 + pkg/cloud/fake/fake.go | 4 + pkg/driver/controller.go | 123 +++++++++++++----- pkg/driver/controller_restore_retry_test.go | 135 ++++++++++++++++++++ 5 files changed, 275 insertions(+), 32 deletions(-) create mode 100644 pkg/cloud/capabilities.go create mode 100644 pkg/driver/controller_restore_retry_test.go 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 56551c7..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) 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/driver/controller.go b/pkg/driver/controller.go index a554a84..562e39e 100644 --- a/pkg/driver/controller.go +++ b/pkg/driver/controller.go @@ -103,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 { @@ -139,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. @@ -214,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) { @@ -281,6 +284,15 @@ func (cs *controllerServer) createVolumeFromSnapshot(ctx context.Context, req *c 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()) @@ -292,17 +304,64 @@ func (cs *controllerServer) createVolumeFromSnapshot(ctx context.Context, req *c 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: volFromSnapshot.ID, - CapacityBytes: volFromSnapshot.Size, + VolumeId: vol.ID, + CapacityBytes: vol.Size, VolumeContext: req.GetParameters(), ContentSource: req.GetVolumeContentSource(), AccessibleTopology: []*csi.Topology{ - Topology{ZoneID: volFromSnapshot.ZoneID}.ToCSI(), + Topology{ZoneID: vol.ZoneID}.ToCSI(), }, }, - }, nil + } } // growVolumeFromSnapshot resizes a volume created from a snapshot when it is 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) + } +}