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/7] 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/7] 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/7] 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/7] 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/7] 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/7] 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/7] 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) + } + }) + } +}