diff --git a/fulfillment-service/internal/cmd/service/start/controller/start_controller_cmd.go b/fulfillment-service/internal/cmd/service/start/controller/start_controller_cmd.go index 3a47ef3e2a..4f9edda619 100644 --- a/fulfillment-service/internal/cmd/service/start/controller/start_controller_cmd.go +++ b/fulfillment-service/internal/cmd/service/start/controller/start_controller_cmd.go @@ -59,6 +59,7 @@ import ( "github.com/osac-project/osac/fulfillment-service/internal/controllers/tenant" "github.com/osac-project/osac/fulfillment-service/internal/controllers/user" "github.com/osac-project/osac/fulfillment-service/internal/controllers/virtualnetwork" + "github.com/osac-project/osac/fulfillment-service/internal/controllers/volume" internalhealth "github.com/osac-project/osac/fulfillment-service/internal/health" "github.com/osac-project/osac/fulfillment-service/internal/idp" hubscheme "github.com/osac-project/osac/fulfillment-service/internal/kubernetes/scheme" @@ -756,6 +757,43 @@ func (r *runnerContext) run(cmd *cobra.Command, argv []string) error { //nolint: } }() + // Create the volume reconciler: + r.logger.InfoContext(ctx, "Creating volume reconciler") + volumeReconcilerFunction, err := volume.NewFunction(). + SetLogger(r.logger). + SetConnection(r.client). + SetHubCache(hubCache). + Build() + if err != nil { + return fmt.Errorf("failed to create volume reconciler function: %w", err) + } + volumeReconciler, err := controllers.NewReconciler[*privatev1.Volume](). + SetLogger(r.logger). + SetName("volume"). + SetClient(r.client). + SetFunction(volumeReconcilerFunction). + SetEventFilter("has(event.volume) || (has(event.hub) && event.type == EVENT_TYPE_OBJECT_CREATED)"). + SetHealthReporter(healthAggregator). + Build() + if err != nil { + return fmt.Errorf("failed to create volume reconciler: %w", err) + } + + // Start the volume reconciler: + r.logger.InfoContext(ctx, "Starting volume reconciler") + go func() { + err := volumeReconciler.Start(ctx) + if err == nil || errors.Is(err, context.Canceled) { + r.logger.InfoContext(ctx, "Volume reconciler finished") + } else { + r.logger.InfoContext( + ctx, + "Volume reconciler failed", + slog.Any("error", err), + ) + } + }() + // Create the role reconciler: r.logger.InfoContext(ctx, "Creating role reconciler") roleReconcilerFunction, err := role.NewFunction(). diff --git a/fulfillment-service/internal/controllers/volume/hubs_client_mock.go b/fulfillment-service/internal/controllers/volume/hubs_client_mock.go new file mode 100644 index 0000000000..f2f7af046b --- /dev/null +++ b/fulfillment-service/internal/controllers/volume/hubs_client_mock.go @@ -0,0 +1,325 @@ +// Code generated by MockGen. DO NOT EDIT. +// Source: ../../api/osac/private/v1/hubs_service_grpc.pb.go +// +// Generated by this command: +// +// mockgen -source=../../api/osac/private/v1/hubs_service_grpc.pb.go -destination=hubs_client_mock.go -package=volume HubsClient +// + +// Package volume is a generated GoMock package. +package volume + +import ( + context "context" + reflect "reflect" + + privatev1 "github.com/osac-project/osac/fulfillment-service/internal/api/osac/private/v1" + gomock "go.uber.org/mock/gomock" + grpc "google.golang.org/grpc" +) + +// MockHubsClient is a mock of HubsClient interface. +type MockHubsClient struct { + ctrl *gomock.Controller + recorder *MockHubsClientMockRecorder + isgomock struct{} +} + +// MockHubsClientMockRecorder is the mock recorder for MockHubsClient. +type MockHubsClientMockRecorder struct { + mock *MockHubsClient +} + +// NewMockHubsClient creates a new mock instance. +func NewMockHubsClient(ctrl *gomock.Controller) *MockHubsClient { + mock := &MockHubsClient{ctrl: ctrl} + mock.recorder = &MockHubsClientMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockHubsClient) EXPECT() *MockHubsClientMockRecorder { + return m.recorder +} + +// Create mocks base method. +func (m *MockHubsClient) Create(ctx context.Context, in *privatev1.HubsCreateRequest, opts ...grpc.CallOption) (*privatev1.HubsCreateResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Create", varargs...) + ret0, _ := ret[0].(*privatev1.HubsCreateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Create indicates an expected call of Create. +func (mr *MockHubsClientMockRecorder) Create(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Create", reflect.TypeOf((*MockHubsClient)(nil).Create), varargs...) +} + +// Delete mocks base method. +func (m *MockHubsClient) Delete(ctx context.Context, in *privatev1.HubsDeleteRequest, opts ...grpc.CallOption) (*privatev1.HubsDeleteResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Delete", varargs...) + ret0, _ := ret[0].(*privatev1.HubsDeleteResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Delete indicates an expected call of Delete. +func (mr *MockHubsClientMockRecorder) Delete(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Delete", reflect.TypeOf((*MockHubsClient)(nil).Delete), varargs...) +} + +// Get mocks base method. +func (m *MockHubsClient) Get(ctx context.Context, in *privatev1.HubsGetRequest, opts ...grpc.CallOption) (*privatev1.HubsGetResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Get", varargs...) + ret0, _ := ret[0].(*privatev1.HubsGetResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Get indicates an expected call of Get. +func (mr *MockHubsClientMockRecorder) Get(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Get", reflect.TypeOf((*MockHubsClient)(nil).Get), varargs...) +} + +// List mocks base method. +func (m *MockHubsClient) List(ctx context.Context, in *privatev1.HubsListRequest, opts ...grpc.CallOption) (*privatev1.HubsListResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "List", varargs...) + ret0, _ := ret[0].(*privatev1.HubsListResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// List indicates an expected call of List. +func (mr *MockHubsClientMockRecorder) List(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "List", reflect.TypeOf((*MockHubsClient)(nil).List), varargs...) +} + +// Signal mocks base method. +func (m *MockHubsClient) Signal(ctx context.Context, in *privatev1.HubsSignalRequest, opts ...grpc.CallOption) (*privatev1.HubsSignalResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Signal", varargs...) + ret0, _ := ret[0].(*privatev1.HubsSignalResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Signal indicates an expected call of Signal. +func (mr *MockHubsClientMockRecorder) Signal(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Signal", reflect.TypeOf((*MockHubsClient)(nil).Signal), varargs...) +} + +// Update mocks base method. +func (m *MockHubsClient) Update(ctx context.Context, in *privatev1.HubsUpdateRequest, opts ...grpc.CallOption) (*privatev1.HubsUpdateResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Update", varargs...) + ret0, _ := ret[0].(*privatev1.HubsUpdateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Update indicates an expected call of Update. +func (mr *MockHubsClientMockRecorder) Update(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Update", reflect.TypeOf((*MockHubsClient)(nil).Update), varargs...) +} + +// MockHubsServer is a mock of HubsServer interface. +type MockHubsServer struct { + ctrl *gomock.Controller + recorder *MockHubsServerMockRecorder + isgomock struct{} +} + +// MockHubsServerMockRecorder is the mock recorder for MockHubsServer. +type MockHubsServerMockRecorder struct { + mock *MockHubsServer +} + +// NewMockHubsServer creates a new mock instance. +func NewMockHubsServer(ctrl *gomock.Controller) *MockHubsServer { + mock := &MockHubsServer{ctrl: ctrl} + mock.recorder = &MockHubsServerMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockHubsServer) EXPECT() *MockHubsServerMockRecorder { + return m.recorder +} + +// Create mocks base method. +func (m *MockHubsServer) Create(arg0 context.Context, arg1 *privatev1.HubsCreateRequest) (*privatev1.HubsCreateResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Create", arg0, arg1) + ret0, _ := ret[0].(*privatev1.HubsCreateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Create indicates an expected call of Create. +func (mr *MockHubsServerMockRecorder) Create(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Create", reflect.TypeOf((*MockHubsServer)(nil).Create), arg0, arg1) +} + +// Delete mocks base method. +func (m *MockHubsServer) Delete(arg0 context.Context, arg1 *privatev1.HubsDeleteRequest) (*privatev1.HubsDeleteResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Delete", arg0, arg1) + ret0, _ := ret[0].(*privatev1.HubsDeleteResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Delete indicates an expected call of Delete. +func (mr *MockHubsServerMockRecorder) Delete(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Delete", reflect.TypeOf((*MockHubsServer)(nil).Delete), arg0, arg1) +} + +// Get mocks base method. +func (m *MockHubsServer) Get(arg0 context.Context, arg1 *privatev1.HubsGetRequest) (*privatev1.HubsGetResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Get", arg0, arg1) + ret0, _ := ret[0].(*privatev1.HubsGetResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Get indicates an expected call of Get. +func (mr *MockHubsServerMockRecorder) Get(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Get", reflect.TypeOf((*MockHubsServer)(nil).Get), arg0, arg1) +} + +// List mocks base method. +func (m *MockHubsServer) List(arg0 context.Context, arg1 *privatev1.HubsListRequest) (*privatev1.HubsListResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "List", arg0, arg1) + ret0, _ := ret[0].(*privatev1.HubsListResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// List indicates an expected call of List. +func (mr *MockHubsServerMockRecorder) List(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "List", reflect.TypeOf((*MockHubsServer)(nil).List), arg0, arg1) +} + +// Signal mocks base method. +func (m *MockHubsServer) Signal(arg0 context.Context, arg1 *privatev1.HubsSignalRequest) (*privatev1.HubsSignalResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Signal", arg0, arg1) + ret0, _ := ret[0].(*privatev1.HubsSignalResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Signal indicates an expected call of Signal. +func (mr *MockHubsServerMockRecorder) Signal(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Signal", reflect.TypeOf((*MockHubsServer)(nil).Signal), arg0, arg1) +} + +// Update mocks base method. +func (m *MockHubsServer) Update(arg0 context.Context, arg1 *privatev1.HubsUpdateRequest) (*privatev1.HubsUpdateResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Update", arg0, arg1) + ret0, _ := ret[0].(*privatev1.HubsUpdateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Update indicates an expected call of Update. +func (mr *MockHubsServerMockRecorder) Update(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Update", reflect.TypeOf((*MockHubsServer)(nil).Update), arg0, arg1) +} + +// mustEmbedUnimplementedHubsServer mocks base method. +func (m *MockHubsServer) mustEmbedUnimplementedHubsServer() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "mustEmbedUnimplementedHubsServer") +} + +// mustEmbedUnimplementedHubsServer indicates an expected call of mustEmbedUnimplementedHubsServer. +func (mr *MockHubsServerMockRecorder) mustEmbedUnimplementedHubsServer() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "mustEmbedUnimplementedHubsServer", reflect.TypeOf((*MockHubsServer)(nil).mustEmbedUnimplementedHubsServer)) +} + +// MockUnsafeHubsServer is a mock of UnsafeHubsServer interface. +type MockUnsafeHubsServer struct { + ctrl *gomock.Controller + recorder *MockUnsafeHubsServerMockRecorder + isgomock struct{} +} + +// MockUnsafeHubsServerMockRecorder is the mock recorder for MockUnsafeHubsServer. +type MockUnsafeHubsServerMockRecorder struct { + mock *MockUnsafeHubsServer +} + +// NewMockUnsafeHubsServer creates a new mock instance. +func NewMockUnsafeHubsServer(ctrl *gomock.Controller) *MockUnsafeHubsServer { + mock := &MockUnsafeHubsServer{ctrl: ctrl} + mock.recorder = &MockUnsafeHubsServerMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockUnsafeHubsServer) EXPECT() *MockUnsafeHubsServerMockRecorder { + return m.recorder +} + +// mustEmbedUnimplementedHubsServer mocks base method. +func (m *MockUnsafeHubsServer) mustEmbedUnimplementedHubsServer() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "mustEmbedUnimplementedHubsServer") +} + +// mustEmbedUnimplementedHubsServer indicates an expected call of mustEmbedUnimplementedHubsServer. +func (mr *MockUnsafeHubsServerMockRecorder) mustEmbedUnimplementedHubsServer() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "mustEmbedUnimplementedHubsServer", reflect.TypeOf((*MockUnsafeHubsServer)(nil).mustEmbedUnimplementedHubsServer)) +} diff --git a/fulfillment-service/internal/controllers/volume/volume_reconciler_function.go b/fulfillment-service/internal/controllers/volume/volume_reconciler_function.go new file mode 100644 index 0000000000..2feaa11347 --- /dev/null +++ b/fulfillment-service/internal/controllers/volume/volume_reconciler_function.go @@ -0,0 +1,430 @@ +/* +Copyright (c) 2026 Red Hat Inc. + +Licensed 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 volume + +//go:generate mockgen -source=../../api/osac/private/v1/volumes_service_grpc.pb.go -destination=volumes_client_mock.go -package=volume VolumesClient +//go:generate mockgen -source=../../api/osac/private/v1/hubs_service_grpc.pb.go -destination=hubs_client_mock.go -package=volume HubsClient + +import ( + "context" + "errors" + "fmt" + "log/slog" + "math/rand/v2" + "slices" + + "google.golang.org/grpc" + "google.golang.org/protobuf/proto" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + clnt "sigs.k8s.io/controller-runtime/pkg/client" + + osacv1alpha1 "github.com/osac-project/osac/osac-operator/api/v1alpha1" + + privatev1 "github.com/osac-project/osac/fulfillment-service/internal/api/osac/private/v1" + "github.com/osac-project/osac/fulfillment-service/internal/controllers" + "github.com/osac-project/osac/fulfillment-service/internal/controllers/finalizers" + "github.com/osac-project/osac/fulfillment-service/internal/kubernetes/annotations" + "github.com/osac-project/osac/fulfillment-service/internal/kubernetes/labels" + "github.com/osac-project/osac/fulfillment-service/internal/masks" +) + +const objectPrefix = "vol-" + +// FunctionBuilder contains the data and logic needed to build a function that +// reconciles volumes. It follows the same builder pattern as other reconciler +// functions (NATGateway, ComputeInstance, etc.). +type FunctionBuilder struct { + logger *slog.Logger + connection *grpc.ClientConn + hubCache controllers.HubCache +} + +// function holds long-lived clients shared across all reconcile invocations. +type function struct { + logger *slog.Logger + hubCache controllers.HubCache + volumesClient privatev1.VolumesClient + hubsClient privatev1.HubsClient + maskCalculator *masks.Calculator +} + +// task carries per-reconciliation state for a single Volume object. +type task struct { + r *function + volume *privatev1.Volume + hubId string + hubNamespace string + hubClient clnt.Client +} + +// NewFunction creates a new builder for the volume reconciler function. +func NewFunction() *FunctionBuilder { + return &FunctionBuilder{} +} + +// SetLogger sets the logger. This is mandatory. +func (b *FunctionBuilder) SetLogger(value *slog.Logger) *FunctionBuilder { + b.logger = value + return b +} + +// SetConnection sets the gRPC client connection. This is mandatory. +func (b *FunctionBuilder) SetConnection(value *grpc.ClientConn) *FunctionBuilder { + b.connection = value + return b +} + +// SetHubCache sets the cache of hubs. This is mandatory. +func (b *FunctionBuilder) SetHubCache(value controllers.HubCache) *FunctionBuilder { + b.hubCache = value + return b +} + +// Build uses the information stored in the builder to create a new volume +// reconciler function. The returned function maps fulfillment-service Volume +// proto objects to Volume CRs on a hub cluster. +func (b *FunctionBuilder) Build() (result controllers.ReconcilerFunction[*privatev1.Volume], err error) { + if b.logger == nil { + err = errors.New("logger is mandatory") + return + } + if b.connection == nil { + err = errors.New("connection is mandatory") + return + } + if b.hubCache == nil { + err = errors.New("hub cache is mandatory") + return + } + + object := &function{ + logger: b.logger, + volumesClient: privatev1.NewVolumesClient(b.connection), + hubsClient: privatev1.NewHubsClient(b.connection), + hubCache: b.hubCache, + maskCalculator: masks.NewCalculator().Build(), + } + result = object.run + return +} + +// run is the ReconcilerFunction entry point. It clones the proto object before +// reconciling, then diffs before/after to compute a FieldMask for a targeted +// gRPC Update that avoids overwriting concurrent changes. +func (r *function) run(ctx context.Context, volume *privatev1.Volume) error { + oldVolume := proto.Clone(volume).(*privatev1.Volume) + t := task{ + r: r, + volume: volume, + } + var err error + if volume.HasMetadata() && volume.GetMetadata().HasDeletionTimestamp() { + err = t.delete(ctx) + } else { + err = t.update(ctx) + } + if err != nil { + return err + } + updateMask := r.maskCalculator.Calculate(oldVolume, volume) + if len(updateMask.GetPaths()) == 0 { + return nil + } + + _, err = r.volumesClient.Update(ctx, privatev1.VolumesUpdateRequest_builder{ + Object: volume, + UpdateMask: updateMask, + }.Build()) + + return err +} + +// update handles the non-delete path: adds the controller finalizer, sets +// default state, validates the tenant, selects a hub, and creates or patches +// the Volume CR on the hub cluster. +func (t *task) update(ctx context.Context) error { + if t.addFinalizer() { + return nil + } + + t.setDefaults() + + if err := t.validateTenant(); err != nil { + return err + } + + if err := t.selectHub(ctx); err != nil { + return err + } + + t.volume.GetStatus().SetHub(t.hubId) + + object, err := t.getKubeObject(ctx) + if err != nil { + return err + } + + spec := t.buildSpec() + + if object == nil { + newObject := &osacv1alpha1.Volume{ + ObjectMeta: metav1.ObjectMeta{ + Namespace: t.hubNamespace, + GenerateName: objectPrefix, + Labels: map[string]string{ + labels.VolumeUuid: t.volume.GetId(), + }, + Annotations: map[string]string{ + annotations.Tenant: t.volume.GetMetadata().GetTenant(), + }, + }, + Spec: spec, + } + err = t.hubClient.Create(ctx, newObject) + if err != nil { + return controllers.HandleK8sWriteError(ctx, t.r.logger, err, t.setFailed) + } + t.r.logger.DebugContext( + ctx, + "Created volume", + slog.String("namespace", newObject.GetNamespace()), + slog.String("name", newObject.GetName()), + ) + } else { + update := object.DeepCopy() + update.Spec = spec + err = t.hubClient.Patch(ctx, update, clnt.MergeFrom(object)) + if err != nil { + return controllers.HandleK8sWriteError(ctx, t.r.logger, err, t.setFailed) + } + t.r.logger.DebugContext( + ctx, + "Updated volume", + slog.String("namespace", object.GetNamespace()), + slog.String("name", object.GetName()), + ) + } + + return nil +} + +func (t *task) setDefaults() { + if !t.volume.HasStatus() { + t.volume.SetStatus(&privatev1.VolumeStatus{}) + } + if t.volume.GetStatus().GetState() == privatev1.VolumeState_VOLUME_STATE_UNSPECIFIED { + t.volume.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_CREATING) + } +} + +func (t *task) validateTenant() error { + if !t.volume.HasMetadata() || t.volume.GetMetadata().GetTenant() == "" { + return errors.New("volume must have a tenant assigned") + } + return nil +} + +// delete handles the deletion path: looks up the hub, finds the Volume CR, +// and deletes it. Removes the controller finalizer when the CR is gone or +// when the hub has been decommissioned. +func (t *task) delete(ctx context.Context) (err error) { + t.hubId = t.volume.GetStatus().GetHub() + if t.hubId == "" { + t.removeFinalizer() + return nil + } + err = t.getHub(ctx) + if err != nil { + if errors.Is(err, controllers.ErrHubNotFound) { + controllers.RemoveFinalizerOnDecommissionedHub(ctx, t.r.logger, t.hubId, "volume_id", t.volume.GetId(), t.removeFinalizer) + return nil + } + return + } + + object, err := t.getKubeObject(ctx) + if err != nil { + return + } + if object == nil { + t.r.logger.DebugContext( + ctx, + "Volume doesn't exist", + slog.String("id", t.volume.GetId()), + ) + t.removeFinalizer() + return + } + + if object.GetDeletionTimestamp() == nil { + err = t.hubClient.Delete(ctx, object) + if err != nil { + return + } + t.r.logger.DebugContext( + ctx, + "Deleted volume", + slog.String("namespace", object.GetNamespace()), + slog.String("name", object.GetName()), + ) + } else { + t.r.logger.DebugContext( + ctx, + "Volume is still being deleted, waiting for K8s finalizers", + slog.String("namespace", object.GetNamespace()), + slog.String("name", object.GetName()), + ) + } + + return +} + +// selectHub assigns a hub to the volume. Unlike resources that inherit their +// hub from a parent (e.g. NATGateway inherits from VirtualNetwork), volumes +// are independent, so the reconciler picks a hub randomly from the available +// hubs (same as ComputeInstance). +func (t *task) selectHub(ctx context.Context) error { + t.hubId = t.volume.GetStatus().GetHub() + if t.hubId == "" { + response, err := t.r.hubsClient.List(ctx, privatev1.HubsListRequest_builder{}.Build()) + if err != nil { + return err + } + if response == nil || len(response.Items) == 0 { + return errors.New("no hubs available") + } + t.hubId = response.Items[rand.IntN(len(response.Items))].GetId() + } + t.r.logger.DebugContext( + ctx, + "Selected hub", + slog.String("id", t.hubId), + ) + hubEntry, err := t.r.hubCache.Get(ctx, t.hubId) + if err != nil { + return err + } + t.hubNamespace = hubEntry.Namespace + t.hubClient = hubEntry.Client + return nil +} + +func (t *task) getHub(ctx context.Context) error { + t.hubId = t.volume.GetStatus().GetHub() + hubEntry, err := t.r.hubCache.Get(ctx, t.hubId) + if err != nil { + return err + } + t.hubNamespace = hubEntry.Namespace + t.hubClient = hubEntry.Client + return nil +} + +// getKubeObject finds the Volume CR on the hub cluster by the UUID label. +// Returns nil if no CR exists yet (first reconcile). +func (t *task) getKubeObject(ctx context.Context) (result *osacv1alpha1.Volume, err error) { + list := &osacv1alpha1.VolumeList{} + err = t.hubClient.List( + ctx, list, + clnt.InNamespace(t.hubNamespace), + clnt.MatchingLabels{ + labels.VolumeUuid: t.volume.GetId(), + }, + ) + if err != nil { + return + } + items := list.Items + count := len(items) + if count > 1 { + err = fmt.Errorf( + "expected at most one volume with identifier '%s' but found %d", + t.volume.GetId(), count, + ) + return + } + if count > 0 { + result = &items[0] + } + return +} + +func (t *task) addFinalizer() bool { + if !t.volume.HasMetadata() { + t.volume.SetMetadata(&privatev1.Metadata{}) + } + list := t.volume.GetMetadata().GetFinalizers() + if !slices.Contains(list, finalizers.Controller) { + list = append(list, finalizers.Controller) + t.volume.GetMetadata().SetFinalizers(list) + return true + } + return false +} + +func (t *task) removeFinalizer() { + if !t.volume.HasMetadata() { + return + } + list := t.volume.GetMetadata().GetFinalizers() + if slices.Contains(list, finalizers.Controller) { + list = slices.DeleteFunc(list, func(item string) bool { + return item == finalizers.Controller + }) + t.volume.GetMetadata().SetFinalizers(list) + } +} + +func (t *task) setFailed(err error) { + if !t.volume.HasStatus() { + t.volume.SetStatus(&privatev1.VolumeStatus{}) + } + t.volume.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_FAILED) + t.volume.GetStatus().SetMessage(err.Error()) +} + +// buildSpec maps the proto VolumeSpec to the osac-operator CRD VolumeSpec. +// Access mode is converted from the proto enum to the CRD typed string. +func (t *task) buildSpec() osacv1alpha1.VolumeSpec { + spec := osacv1alpha1.VolumeSpec{ + StorageTier: t.volume.GetSpec().GetStorageTier(), + SizeGiB: t.volume.GetSpec().GetSizeGib(), + AccessMode: protoAccessModeToCRD(t.volume.GetSpec().GetAccessMode()), + } + if t.volume.GetSpec().HasPvcRef() { + spec.PVCRef = &osacv1alpha1.PVCReferenceType{ + Name: t.volume.GetSpec().GetPvcRef().GetName(), + Namespace: t.volume.GetSpec().GetPvcRef().GetNamespace(), + Cluster: t.volume.GetSpec().GetPvcRef().GetCluster(), + } + } + return spec +} + +// protoAccessModeToCRD converts the proto VolumeAccessMode enum to the CRD +// VolumeAccessMode typed string. +func protoAccessModeToCRD(mode privatev1.VolumeAccessMode) osacv1alpha1.VolumeAccessMode { + switch mode { + case privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_ONCE: + return osacv1alpha1.VolumeAccessModeReadWriteOnce + case privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_ONLY_MANY: + return osacv1alpha1.VolumeAccessModeReadOnlyMany + case privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_MANY: + return osacv1alpha1.VolumeAccessModeReadWriteMany + case privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_ONCE_POD: + return osacv1alpha1.VolumeAccessModeReadWriteOncePod + default: + return osacv1alpha1.VolumeAccessModeReadWriteOnce + } +} diff --git a/fulfillment-service/internal/controllers/volume/volume_reconciler_function_test.go b/fulfillment-service/internal/controllers/volume/volume_reconciler_function_test.go new file mode 100644 index 0000000000..2083d65f8d --- /dev/null +++ b/fulfillment-service/internal/controllers/volume/volume_reconciler_function_test.go @@ -0,0 +1,716 @@ +/* +Copyright (c) 2026 Red Hat Inc. + +Licensed 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 volume + +import ( + "context" + "errors" + "slices" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + "go.uber.org/mock/gomock" + "google.golang.org/grpc" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/util/validation/field" + clnt "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" + + privatev1 "github.com/osac-project/osac/fulfillment-service/internal/api/osac/private/v1" + "github.com/osac-project/osac/fulfillment-service/internal/controllers" + "github.com/osac-project/osac/fulfillment-service/internal/controllers/finalizers" + "github.com/osac-project/osac/fulfillment-service/internal/kubernetes/labels" + "github.com/osac-project/osac/fulfillment-service/internal/masks" + osacv1alpha1 "github.com/osac-project/osac/osac-operator/api/v1alpha1" +) + +func newVolumeCR(id, namespace, name string, deletionTimestamp *metav1.Time) *osacv1alpha1.Volume { + obj := &osacv1alpha1.Volume{ + ObjectMeta: metav1.ObjectMeta{ + Namespace: namespace, + Name: name, + Labels: map[string]string{ + labels.VolumeUuid: id, + }, + }, + } + if deletionTimestamp != nil { + obj.SetDeletionTimestamp(deletionTimestamp) + obj.SetFinalizers([]string{"osac.openshift.io/volume"}) + } + return obj +} + +func hasFinalizer(volume *privatev1.Volume) bool { + return slices.Contains(volume.GetMetadata().GetFinalizers(), finalizers.Controller) +} + +func newTaskForDelete(volumeID, hubID string, hubCache controllers.HubCache) *task { + volume := privatev1.Volume_builder{ + Id: volumeID, + Metadata: privatev1.Metadata_builder{ + Finalizers: []string{finalizers.Controller}, + }.Build(), + Status: privatev1.VolumeStatus_builder{ + Hub: hubID, + }.Build(), + }.Build() + + f := &function{ + logger: logger, + hubCache: hubCache, + } + + return &task{ + r: f, + volume: volume, + } +} + +var _ = Describe("buildSpec", func() { + It("maps all spec fields including access mode enum", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-buildspec-1", + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "gold", + SizeGib: 100, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_ONCE, + }.Build(), + }.Build(), + } + + spec := t.buildSpec() + + Expect(spec.StorageTier).To(Equal("gold")) + Expect(spec.SizeGiB).To(Equal(int64(100))) + Expect(spec.AccessMode).To(Equal(osacv1alpha1.VolumeAccessModeReadWriteOnce)) + Expect(spec.PVCRef).To(BeNil()) + }) + + It("maps ReadWriteMany access mode", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-buildspec-rwm", + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "silver", + SizeGib: 50, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_MANY, + }.Build(), + }.Build(), + } + + spec := t.buildSpec() + Expect(spec.AccessMode).To(Equal(osacv1alpha1.VolumeAccessModeReadWriteMany)) + }) + + It("maps ReadWriteOncePod access mode", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-buildspec-rwop", + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "silver", + SizeGib: 50, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_ONCE_POD, + }.Build(), + }.Build(), + } + + spec := t.buildSpec() + Expect(spec.AccessMode).To(Equal(osacv1alpha1.VolumeAccessModeReadWriteOncePod)) + }) + + It("maps ReadOnlyMany access mode", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-buildspec-rom", + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "silver", + SizeGib: 50, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_ONLY_MANY, + }.Build(), + }.Build(), + } + + spec := t.buildSpec() + Expect(spec.AccessMode).To(Equal(osacv1alpha1.VolumeAccessModeReadOnlyMany)) + }) + + It("includes PVC reference when present", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-buildspec-pvc", + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "gold", + SizeGib: 100, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_ONCE, + PvcRef: privatev1.PVCReference_builder{ + Name: "my-pvc", + Namespace: "my-ns", + Cluster: "cluster-1", + }.Build(), + }.Build(), + }.Build(), + } + + spec := t.buildSpec() + + Expect(spec.PVCRef).ToNot(BeNil()) + Expect(spec.PVCRef.Name).To(Equal("my-pvc")) + Expect(spec.PVCRef.Namespace).To(Equal("my-ns")) + Expect(spec.PVCRef.Cluster).To(Equal("cluster-1")) + }) + + It("defaults unspecified access mode to ReadWriteOnce", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-buildspec-default", + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "gold", + SizeGib: 100, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_UNSPECIFIED, + }.Build(), + }.Build(), + } + + spec := t.buildSpec() + Expect(spec.AccessMode).To(Equal(osacv1alpha1.VolumeAccessModeReadWriteOnce)) + }) +}) + +var _ = Describe("setDefaults", func() { + It("sets CREATING state when status is unspecified", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-defaults-1", + }.Build(), + } + + t.setDefaults() + + Expect(t.volume.GetStatus().GetState()).To( + Equal(privatev1.VolumeState_VOLUME_STATE_CREATING), + ) + }) + + It("does not overwrite existing state", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-defaults-existing", + Status: privatev1.VolumeStatus_builder{ + State: privatev1.VolumeState_VOLUME_STATE_AVAILABLE, + }.Build(), + }.Build(), + } + + t.setDefaults() + + Expect(t.volume.GetStatus().GetState()).To( + Equal(privatev1.VolumeState_VOLUME_STATE_AVAILABLE), + ) + }) + + It("creates status if it doesn't exist", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-defaults-no-status", + }.Build(), + } + + Expect(t.volume.HasStatus()).To(BeFalse()) + + t.setDefaults() + + Expect(t.volume.HasStatus()).To(BeTrue()) + Expect(t.volume.GetStatus().GetState()).To( + Equal(privatev1.VolumeState_VOLUME_STATE_CREATING), + ) + }) +}) + +var _ = Describe("validateTenant", func() { + It("succeeds when a tenant is assigned", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Metadata: privatev1.Metadata_builder{ + Tenant: "tenant-1", + }.Build(), + }.Build(), + } + + err := t.validateTenant() + Expect(err).ToNot(HaveOccurred()) + }) + + It("fails when tenant is empty", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Metadata: privatev1.Metadata_builder{ + Tenant: "", + }.Build(), + }.Build(), + } + + err := t.validateTenant() + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("tenant")) + }) + + It("fails when metadata is missing", func() { + t := &task{ + volume: privatev1.Volume_builder{}.Build(), + } + + err := t.validateTenant() + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("tenant")) + }) +}) + +var _ = Describe("addFinalizer", func() { + It("adds finalizer when not present", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-add-fin", + Metadata: privatev1.Metadata_builder{ + Finalizers: []string{}, + }.Build(), + }.Build(), + } + + added := t.addFinalizer() + + Expect(added).To(BeTrue()) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + }) + + It("does not add finalizer when already present", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-has-fin", + Metadata: privatev1.Metadata_builder{ + Finalizers: []string{finalizers.Controller}, + }.Build(), + }.Build(), + } + + added := t.addFinalizer() + + Expect(added).To(BeFalse()) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + Expect(t.volume.GetMetadata().GetFinalizers()).To(HaveLen(1)) + }) + + It("creates metadata if it doesn't exist", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-no-meta", + }.Build(), + } + + Expect(t.volume.HasMetadata()).To(BeFalse()) + + added := t.addFinalizer() + + Expect(added).To(BeTrue()) + Expect(t.volume.HasMetadata()).To(BeTrue()) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + }) +}) + +var _ = Describe("removeFinalizer", func() { + It("removes finalizer when present", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-rm-fin", + Metadata: privatev1.Metadata_builder{ + Finalizers: []string{finalizers.Controller, "other-finalizer"}, + }.Build(), + }.Build(), + } + + Expect(hasFinalizer(t.volume)).To(BeTrue()) + + t.removeFinalizer() + + Expect(hasFinalizer(t.volume)).To(BeFalse()) + Expect(t.volume.GetMetadata().GetFinalizers()).To(ContainElement("other-finalizer")) + }) + + It("does nothing when finalizer not present", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-no-fin", + Metadata: privatev1.Metadata_builder{ + Finalizers: []string{"other-finalizer"}, + }.Build(), + }.Build(), + } + + t.removeFinalizer() + + Expect(hasFinalizer(t.volume)).To(BeFalse()) + Expect(t.volume.GetMetadata().GetFinalizers()).To(ContainElement("other-finalizer")) + }) + + It("does nothing when metadata doesn't exist", func() { + t := &task{ + volume: privatev1.Volume_builder{ + Id: "vol-no-meta-rm", + }.Build(), + } + + t.removeFinalizer() + + Expect(t.volume.HasMetadata()).To(BeFalse()) + }) +}) + +var _ = Describe("delete", func() { + const ( + volumeID = "vol-delete-id" + hubID = "test-hub" + hubNamespace = "test-ns" + crName = "vol-test" + ) + + var ( + ctx context.Context + ctrl *gomock.Controller + ) + + BeforeEach(func() { + ctx = context.Background() + ctrl = gomock.NewController(GinkgoT()) + DeferCleanup(ctrl.Finish) + }) + + It("removes finalizer when no hub is assigned", func() { + volume := privatev1.Volume_builder{ + Id: volumeID, + Metadata: privatev1.Metadata_builder{ + Finalizers: []string{finalizers.Controller}, + }.Build(), + Status: privatev1.VolumeStatus_builder{}.Build(), + }.Build() + + f := &function{logger: logger} + t := &task{r: f, volume: volume} + + Expect(hasFinalizer(t.volume)).To(BeTrue()) + + err := t.delete(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(hasFinalizer(t.volume)).To(BeFalse()) + }) + + It("removes finalizer when K8s object doesn't exist", func() { + scheme := runtime.NewScheme() + Expect(osacv1alpha1.AddToScheme(scheme)).To(Succeed()) + fakeClient := fake.NewClientBuilder(). + WithScheme(scheme). + Build() + + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), hubID). + Return(&controllers.HubEntry{ + Namespace: hubNamespace, + Client: fakeClient, + }, nil) + + t := newTaskForDelete(volumeID, hubID, hubCache) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + + err := t.delete(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(hasFinalizer(t.volume)).To(BeFalse()) + }) + + It("calls hubClient.Delete when K8s object exists without DeletionTimestamp", func() { + cr := newVolumeCR(volumeID, hubNamespace, crName, nil) + + scheme := runtime.NewScheme() + Expect(osacv1alpha1.AddToScheme(scheme)).To(Succeed()) + + deleteCalled := false + fakeClient := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(cr). + WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(ctx context.Context, client clnt.WithWatch, obj clnt.Object, opts ...clnt.DeleteOption) error { + deleteCalled = true + return nil + }, + }). + Build() + + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), hubID). + Return(&controllers.HubEntry{ + Namespace: hubNamespace, + Client: fakeClient, + }, nil) + + t := newTaskForDelete(volumeID, hubID, hubCache) + + err := t.delete(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(deleteCalled).To(BeTrue()) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + }) + + It("does not call hubClient.Delete when K8s object has DeletionTimestamp", func() { + now := metav1.Now() + cr := newVolumeCR(volumeID, hubNamespace, crName, &now) + + scheme := runtime.NewScheme() + Expect(osacv1alpha1.AddToScheme(scheme)).To(Succeed()) + + deleteCalled := false + fakeClient := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(cr). + WithInterceptorFuncs(interceptor.Funcs{ + Delete: func(ctx context.Context, client clnt.WithWatch, obj clnt.Object, opts ...clnt.DeleteOption) error { + deleteCalled = true + return nil + }, + }). + Build() + + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), hubID). + Return(&controllers.HubEntry{ + Namespace: hubNamespace, + Client: fakeClient, + }, nil) + + t := newTaskForDelete(volumeID, hubID, hubCache) + + err := t.delete(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(deleteCalled).To(BeFalse()) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + }) + + It("removes finalizer when hub is decommissioned", func() { + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), hubID). + Return(nil, controllers.ErrHubNotFound) + + t := newTaskForDelete(volumeID, hubID, hubCache) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + + err := t.delete(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(hasFinalizer(t.volume)).To(BeFalse()) + }) + + It("propagates error when hub cache returns transient error", func() { + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), hubID). + Return(nil, errors.New("connection refused")) + + t := newTaskForDelete(volumeID, hubID, hubCache) + + err := t.delete(ctx) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("connection refused")) + Expect(hasFinalizer(t.volume)).To(BeTrue()) + }) +}) + +var _ = Describe("selectHub", func() { + var ( + ctx context.Context + ctrl *gomock.Controller + ) + + BeforeEach(func() { + ctx = context.Background() + ctrl = gomock.NewController(GinkgoT()) + DeferCleanup(ctrl.Finish) + }) + + It("uses existing hub from status", func() { + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), "hub-1"). + Return(&controllers.HubEntry{ + Namespace: "hub-ns", + Client: fake.NewClientBuilder().Build(), + }, nil) + + t := &task{ + r: &function{ + logger: logger, + hubCache: hubCache, + }, + volume: privatev1.Volume_builder{ + Id: "vol-existing-hub", + Status: privatev1.VolumeStatus_builder{ + Hub: "hub-1", + }.Build(), + }.Build(), + } + + err := t.selectHub(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(t.hubId).To(Equal("hub-1")) + Expect(t.hubNamespace).To(Equal("hub-ns")) + }) + + It("selects a hub randomly when status hub is empty", func() { + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), gomock.Any()). + Return(&controllers.HubEntry{ + Namespace: "selected-ns", + Client: fake.NewClientBuilder().Build(), + }, nil) + + hubsClient := NewMockHubsClient(ctrl) + hubsClient.EXPECT(). + List(gomock.Any(), gomock.Any(), gomock.Any()). + Return(privatev1.HubsListResponse_builder{ + Items: []*privatev1.Hub{ + privatev1.Hub_builder{Id: "hub-a"}.Build(), + }, + }.Build(), nil) + + t := &task{ + r: &function{ + logger: logger, + hubCache: hubCache, + hubsClient: hubsClient, + }, + volume: privatev1.Volume_builder{ + Id: "vol-no-hub", + }.Build(), + } + + err := t.selectHub(ctx) + Expect(err).ToNot(HaveOccurred()) + Expect(t.hubId).To(Equal("hub-a")) + Expect(t.hubNamespace).To(Equal("selected-ns")) + }) + + It("returns error when no hubs are available", func() { + hubsClient := NewMockHubsClient(ctrl) + hubsClient.EXPECT(). + List(gomock.Any(), gomock.Any(), gomock.Any()). + Return(privatev1.HubsListResponse_builder{ + Items: []*privatev1.Hub{}, + }.Build(), nil) + + t := &task{ + r: &function{ + logger: logger, + hubsClient: hubsClient, + }, + volume: privatev1.Volume_builder{ + Id: "vol-no-hubs", + }.Build(), + } + + err := t.selectHub(ctx) + Expect(err).To(HaveOccurred()) + Expect(err.Error()).To(ContainSubstring("no hubs available")) + }) +}) + +var _ = Describe("Kubernetes validation error handling", func() { + It("sets state to FAILED when K8s Create returns Invalid error", func() { + ctx := context.Background() + ctrl := gomock.NewController(GinkgoT()) + DeferCleanup(ctrl.Finish) + + scheme := runtime.NewScheme() + Expect(osacv1alpha1.AddToScheme(scheme)).To(Succeed()) + + fakeClient := fake.NewClientBuilder(). + WithScheme(scheme). + WithInterceptorFuncs(interceptor.Funcs{ + Create: func(ctx context.Context, client clnt.WithWatch, obj clnt.Object, opts ...clnt.CreateOption) error { + return apierrors.NewInvalid( + schema.GroupKind{Group: "osac.openshift.io", Kind: "Volume"}, + "vol-test", + field.ErrorList{ + field.Invalid( + field.NewPath("spec", "storageTier"), + "invalid-value", + "spec.storageTier is invalid", + ), + }, + ) + }, + }). + Build() + + hubCache := controllers.NewMockHubCache(ctrl) + hubCache.EXPECT(). + Get(gomock.Any(), "hub-1"). + Return(&controllers.HubEntry{Namespace: "test-ns", Client: fakeClient}, nil). + AnyTimes() + + volumesClient := NewMockVolumesClient(ctrl) + volumesClient.EXPECT(). + Update(gomock.Any(), gomock.Any(), gomock.Any()). + DoAndReturn(func(ctx context.Context, req *privatev1.VolumesUpdateRequest, opts ...grpc.CallOption) (*privatev1.VolumesUpdateResponse, error) { + return &privatev1.VolumesUpdateResponse{Object: req.GetObject()}, nil + }). + MinTimes(1) + + volume := privatev1.Volume_builder{ + Id: "vol-validation-test", + Metadata: privatev1.Metadata_builder{ + Finalizers: []string{finalizers.Controller}, + Tenant: "test-tenant", + }.Build(), + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "gold", + SizeGib: 100, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_ONCE, + }.Build(), + Status: privatev1.VolumeStatus_builder{ + State: privatev1.VolumeState_VOLUME_STATE_CREATING, + Hub: "hub-1", + }.Build(), + }.Build() + + f := &function{ + logger: logger, + hubCache: hubCache, + volumesClient: volumesClient, + maskCalculator: masks.NewCalculator().Build(), + } + + err := f.run(ctx, volume) + Expect(err).ToNot(HaveOccurred()) + + Expect(volume.GetStatus().GetState()).To( + Equal(privatev1.VolumeState_VOLUME_STATE_FAILED), + ) + Expect(volume.GetStatus().GetMessage()).To(ContainSubstring("spec.storageTier")) + }) +}) diff --git a/fulfillment-service/internal/controllers/volume/volume_suite_test.go b/fulfillment-service/internal/controllers/volume/volume_suite_test.go new file mode 100644 index 0000000000..079cb914f3 --- /dev/null +++ b/fulfillment-service/internal/controllers/volume/volume_suite_test.go @@ -0,0 +1,39 @@ +/* +Copyright (c) 2026 Red Hat Inc. + +Licensed 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 volume + +import ( + "log/slog" + "testing" + + . "github.com/onsi/ginkgo/v2/dsl/core" + . "github.com/onsi/gomega" + "github.com/osac-project/osac/fulfillment-service/internal/logging" +) + +func TestVolume(t *testing.T) { + RegisterFailHandler(Fail) + RunSpecs(t, "Volume controller") +} + +var logger *slog.Logger + +var _ = BeforeSuite(func() { + var err error + logger, err = logging.NewLogger(). + SetLevel(slog.LevelDebug.String()). + SetWriter(GinkgoWriter). + Build() + Expect(err).ToNot(HaveOccurred()) +}) diff --git a/fulfillment-service/internal/controllers/volume/volumes_client_mock.go b/fulfillment-service/internal/controllers/volume/volumes_client_mock.go new file mode 100644 index 0000000000..e9c9fa639b --- /dev/null +++ b/fulfillment-service/internal/controllers/volume/volumes_client_mock.go @@ -0,0 +1,325 @@ +// Code generated by MockGen. DO NOT EDIT. +// Source: ../../api/osac/private/v1/volumes_service_grpc.pb.go +// +// Generated by this command: +// +// mockgen -source=../../api/osac/private/v1/volumes_service_grpc.pb.go -destination=volumes_client_mock.go -package=volume VolumesClient +// + +// Package volume is a generated GoMock package. +package volume + +import ( + context "context" + reflect "reflect" + + privatev1 "github.com/osac-project/osac/fulfillment-service/internal/api/osac/private/v1" + gomock "go.uber.org/mock/gomock" + grpc "google.golang.org/grpc" +) + +// MockVolumesClient is a mock of VolumesClient interface. +type MockVolumesClient struct { + ctrl *gomock.Controller + recorder *MockVolumesClientMockRecorder + isgomock struct{} +} + +// MockVolumesClientMockRecorder is the mock recorder for MockVolumesClient. +type MockVolumesClientMockRecorder struct { + mock *MockVolumesClient +} + +// NewMockVolumesClient creates a new mock instance. +func NewMockVolumesClient(ctrl *gomock.Controller) *MockVolumesClient { + mock := &MockVolumesClient{ctrl: ctrl} + mock.recorder = &MockVolumesClientMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockVolumesClient) EXPECT() *MockVolumesClientMockRecorder { + return m.recorder +} + +// Create mocks base method. +func (m *MockVolumesClient) Create(ctx context.Context, in *privatev1.VolumesCreateRequest, opts ...grpc.CallOption) (*privatev1.VolumesCreateResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Create", varargs...) + ret0, _ := ret[0].(*privatev1.VolumesCreateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Create indicates an expected call of Create. +func (mr *MockVolumesClientMockRecorder) Create(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Create", reflect.TypeOf((*MockVolumesClient)(nil).Create), varargs...) +} + +// Delete mocks base method. +func (m *MockVolumesClient) Delete(ctx context.Context, in *privatev1.VolumesDeleteRequest, opts ...grpc.CallOption) (*privatev1.VolumesDeleteResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Delete", varargs...) + ret0, _ := ret[0].(*privatev1.VolumesDeleteResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Delete indicates an expected call of Delete. +func (mr *MockVolumesClientMockRecorder) Delete(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Delete", reflect.TypeOf((*MockVolumesClient)(nil).Delete), varargs...) +} + +// Get mocks base method. +func (m *MockVolumesClient) Get(ctx context.Context, in *privatev1.VolumesGetRequest, opts ...grpc.CallOption) (*privatev1.VolumesGetResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Get", varargs...) + ret0, _ := ret[0].(*privatev1.VolumesGetResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Get indicates an expected call of Get. +func (mr *MockVolumesClientMockRecorder) Get(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Get", reflect.TypeOf((*MockVolumesClient)(nil).Get), varargs...) +} + +// List mocks base method. +func (m *MockVolumesClient) List(ctx context.Context, in *privatev1.VolumesListRequest, opts ...grpc.CallOption) (*privatev1.VolumesListResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "List", varargs...) + ret0, _ := ret[0].(*privatev1.VolumesListResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// List indicates an expected call of List. +func (mr *MockVolumesClientMockRecorder) List(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "List", reflect.TypeOf((*MockVolumesClient)(nil).List), varargs...) +} + +// Signal mocks base method. +func (m *MockVolumesClient) Signal(ctx context.Context, in *privatev1.VolumesSignalRequest, opts ...grpc.CallOption) (*privatev1.VolumesSignalResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Signal", varargs...) + ret0, _ := ret[0].(*privatev1.VolumesSignalResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Signal indicates an expected call of Signal. +func (mr *MockVolumesClientMockRecorder) Signal(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Signal", reflect.TypeOf((*MockVolumesClient)(nil).Signal), varargs...) +} + +// Update mocks base method. +func (m *MockVolumesClient) Update(ctx context.Context, in *privatev1.VolumesUpdateRequest, opts ...grpc.CallOption) (*privatev1.VolumesUpdateResponse, error) { + m.ctrl.T.Helper() + varargs := []any{ctx, in} + for _, a := range opts { + varargs = append(varargs, a) + } + ret := m.ctrl.Call(m, "Update", varargs...) + ret0, _ := ret[0].(*privatev1.VolumesUpdateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Update indicates an expected call of Update. +func (mr *MockVolumesClientMockRecorder) Update(ctx, in any, opts ...any) *gomock.Call { + mr.mock.ctrl.T.Helper() + varargs := append([]any{ctx, in}, opts...) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Update", reflect.TypeOf((*MockVolumesClient)(nil).Update), varargs...) +} + +// MockVolumesServer is a mock of VolumesServer interface. +type MockVolumesServer struct { + ctrl *gomock.Controller + recorder *MockVolumesServerMockRecorder + isgomock struct{} +} + +// MockVolumesServerMockRecorder is the mock recorder for MockVolumesServer. +type MockVolumesServerMockRecorder struct { + mock *MockVolumesServer +} + +// NewMockVolumesServer creates a new mock instance. +func NewMockVolumesServer(ctrl *gomock.Controller) *MockVolumesServer { + mock := &MockVolumesServer{ctrl: ctrl} + mock.recorder = &MockVolumesServerMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockVolumesServer) EXPECT() *MockVolumesServerMockRecorder { + return m.recorder +} + +// Create mocks base method. +func (m *MockVolumesServer) Create(arg0 context.Context, arg1 *privatev1.VolumesCreateRequest) (*privatev1.VolumesCreateResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Create", arg0, arg1) + ret0, _ := ret[0].(*privatev1.VolumesCreateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Create indicates an expected call of Create. +func (mr *MockVolumesServerMockRecorder) Create(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Create", reflect.TypeOf((*MockVolumesServer)(nil).Create), arg0, arg1) +} + +// Delete mocks base method. +func (m *MockVolumesServer) Delete(arg0 context.Context, arg1 *privatev1.VolumesDeleteRequest) (*privatev1.VolumesDeleteResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Delete", arg0, arg1) + ret0, _ := ret[0].(*privatev1.VolumesDeleteResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Delete indicates an expected call of Delete. +func (mr *MockVolumesServerMockRecorder) Delete(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Delete", reflect.TypeOf((*MockVolumesServer)(nil).Delete), arg0, arg1) +} + +// Get mocks base method. +func (m *MockVolumesServer) Get(arg0 context.Context, arg1 *privatev1.VolumesGetRequest) (*privatev1.VolumesGetResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Get", arg0, arg1) + ret0, _ := ret[0].(*privatev1.VolumesGetResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Get indicates an expected call of Get. +func (mr *MockVolumesServerMockRecorder) Get(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Get", reflect.TypeOf((*MockVolumesServer)(nil).Get), arg0, arg1) +} + +// List mocks base method. +func (m *MockVolumesServer) List(arg0 context.Context, arg1 *privatev1.VolumesListRequest) (*privatev1.VolumesListResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "List", arg0, arg1) + ret0, _ := ret[0].(*privatev1.VolumesListResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// List indicates an expected call of List. +func (mr *MockVolumesServerMockRecorder) List(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "List", reflect.TypeOf((*MockVolumesServer)(nil).List), arg0, arg1) +} + +// Signal mocks base method. +func (m *MockVolumesServer) Signal(arg0 context.Context, arg1 *privatev1.VolumesSignalRequest) (*privatev1.VolumesSignalResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Signal", arg0, arg1) + ret0, _ := ret[0].(*privatev1.VolumesSignalResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Signal indicates an expected call of Signal. +func (mr *MockVolumesServerMockRecorder) Signal(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Signal", reflect.TypeOf((*MockVolumesServer)(nil).Signal), arg0, arg1) +} + +// Update mocks base method. +func (m *MockVolumesServer) Update(arg0 context.Context, arg1 *privatev1.VolumesUpdateRequest) (*privatev1.VolumesUpdateResponse, error) { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "Update", arg0, arg1) + ret0, _ := ret[0].(*privatev1.VolumesUpdateResponse) + ret1, _ := ret[1].(error) + return ret0, ret1 +} + +// Update indicates an expected call of Update. +func (mr *MockVolumesServerMockRecorder) Update(arg0, arg1 any) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Update", reflect.TypeOf((*MockVolumesServer)(nil).Update), arg0, arg1) +} + +// mustEmbedUnimplementedVolumesServer mocks base method. +func (m *MockVolumesServer) mustEmbedUnimplementedVolumesServer() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "mustEmbedUnimplementedVolumesServer") +} + +// mustEmbedUnimplementedVolumesServer indicates an expected call of mustEmbedUnimplementedVolumesServer. +func (mr *MockVolumesServerMockRecorder) mustEmbedUnimplementedVolumesServer() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "mustEmbedUnimplementedVolumesServer", reflect.TypeOf((*MockVolumesServer)(nil).mustEmbedUnimplementedVolumesServer)) +} + +// MockUnsafeVolumesServer is a mock of UnsafeVolumesServer interface. +type MockUnsafeVolumesServer struct { + ctrl *gomock.Controller + recorder *MockUnsafeVolumesServerMockRecorder + isgomock struct{} +} + +// MockUnsafeVolumesServerMockRecorder is the mock recorder for MockUnsafeVolumesServer. +type MockUnsafeVolumesServerMockRecorder struct { + mock *MockUnsafeVolumesServer +} + +// NewMockUnsafeVolumesServer creates a new mock instance. +func NewMockUnsafeVolumesServer(ctrl *gomock.Controller) *MockUnsafeVolumesServer { + mock := &MockUnsafeVolumesServer{ctrl: ctrl} + mock.recorder = &MockUnsafeVolumesServerMockRecorder{mock} + return mock +} + +// EXPECT returns an object that allows the caller to indicate expected use. +func (m *MockUnsafeVolumesServer) EXPECT() *MockUnsafeVolumesServerMockRecorder { + return m.recorder +} + +// mustEmbedUnimplementedVolumesServer mocks base method. +func (m *MockUnsafeVolumesServer) mustEmbedUnimplementedVolumesServer() { + m.ctrl.T.Helper() + m.ctrl.Call(m, "mustEmbedUnimplementedVolumesServer") +} + +// mustEmbedUnimplementedVolumesServer indicates an expected call of mustEmbedUnimplementedVolumesServer. +func (mr *MockUnsafeVolumesServerMockRecorder) mustEmbedUnimplementedVolumesServer() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "mustEmbedUnimplementedVolumesServer", reflect.TypeOf((*MockUnsafeVolumesServer)(nil).mustEmbedUnimplementedVolumesServer)) +} diff --git a/fulfillment-service/internal/kubernetes/labels/kubernetes_labels.go b/fulfillment-service/internal/kubernetes/labels/kubernetes_labels.go index 1f51386b78..936a78b301 100644 --- a/fulfillment-service/internal/kubernetes/labels/kubernetes_labels.go +++ b/fulfillment-service/internal/kubernetes/labels/kubernetes_labels.go @@ -60,6 +60,9 @@ var NATGatewayUuid = fmt.Sprintf("%s/%s", group, "natgateway-uuid") // TenantUuid is the label where the fulfillment API will write the identifier of the tenant. var TenantUuid = fmt.Sprintf("%s/%s", group, "tenant-uuid") +// VolumeUuid is the label where the fulfillment API will write the identifier of the volume. +var VolumeUuid = fmt.Sprintf("%s/%s", group, "volume-uuid") + // TenantRef is the label used to reference the tenant object from associated resources (e.g., namespaces). var TenantRef = fmt.Sprintf("%s/%s", group, "tenant-ref") diff --git a/osac-operator/AGENTS.md b/osac-operator/AGENTS.md index a4c42d8c25..563289fc8e 100644 --- a/osac-operator/AGENTS.md +++ b/osac-operator/AGENTS.md @@ -16,6 +16,7 @@ OSAC operator is a Kubernetes operator that reconciles infrastructure resources - **ExternalIP** — external IP allocated from ExternalIPPool - **ExternalIPAttachment** — attachment of ExternalIP to ComputeInstance - **NATGateway** — outbound SNAT for a VirtualNetwork +- **Volume** (`vol`) — block storage on vendor arrays via CSI ## Critical Rules diff --git a/osac-operator/README.md b/osac-operator/README.md index 78eff1a55e..b714b60207 100644 --- a/osac-operator/README.md +++ b/osac-operator/README.md @@ -21,6 +21,8 @@ custom resources and reconciles them to their desired state: - **Subnet** (`subnet`) — represents a subnet within a VirtualNetwork. - **SecurityGroup** (`sg`) — defines network security (firewall) rules with ingress/egress rules, protocols, port ranges, and CIDR blocks. +- **Volume** (`vol`) — provisions block storage on vendor arrays via the + VendorProvisioner interface (vendor CSI controllers). ## Configuration diff --git a/osac-operator/api/v1alpha1/volume_names.go b/osac-operator/api/v1alpha1/volume_names.go deleted file mode 100644 index 9b773388de..0000000000 --- a/osac-operator/api/v1alpha1/volume_names.go +++ /dev/null @@ -1,39 +0,0 @@ -/* -Copyright 2026. - -Licensed 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 v1alpha1 - -const ( - // VolumeNamespace is the default namespace where Volume CRs are created. - VolumeNamespace = "osac-volume" - - // VolumeLabelName is the label key for the volume name. - VolumeLabelName = "osac.openshift.io/volume" - - // VolumeLabelUUID is the label key for the fulfillment-service volume ID, - // used by the feedback controller to map Volume CRs back to inventory records. - VolumeLabelUUID = "osac.openshift.io/volume-uuid" - - // VolumeFinalizer is the finalizer managed by the Volume resource controller. - VolumeFinalizer = "osac.openshift.io/volume-finalizer" - - // VolumeFeedbackFinalizer is the finalizer managed by the Volume feedback controller. - VolumeFeedbackFinalizer = "osac.openshift.io/volume-feedback" - - // VolumeCleanupFinalizer is the finalizer added to ClusterOrder when volumes - // are created, blocking cluster deletion until volumes are processed. - VolumeCleanupFinalizer = "osac.openshift.io/volume-cleanup" -) diff --git a/osac-operator/charts/operator/templates/clusterrole.yaml b/osac-operator/charts/operator/templates/clusterrole.yaml index b960a3b541..fdb8200a67 100644 --- a/osac-operator/charts/operator/templates/clusterrole.yaml +++ b/osac-operator/charts/operator/templates/clusterrole.yaml @@ -92,6 +92,7 @@ rules: - subnets - tenants - virtualnetworks + - volumes verbs: - create - delete @@ -114,6 +115,7 @@ rules: - subnets/finalizers - tenants/finalizers - virtualnetworks/finalizers + - volumes/finalizers verbs: - update - apiGroups: @@ -130,6 +132,7 @@ rules: - subnets/status - tenants/status - virtualnetworks/status + - volumes/status verbs: - get - patch diff --git a/osac-operator/charts/operator/templates/deployment.yaml b/osac-operator/charts/operator/templates/deployment.yaml index b8bb06cf19..d066659011 100644 --- a/osac-operator/charts/operator/templates/deployment.yaml +++ b/osac-operator/charts/operator/templates/deployment.yaml @@ -94,6 +94,10 @@ spec: valueFrom: fieldRef: fieldPath: metadata.namespace + - name: OSAC_VOLUME_NAMESPACE + valueFrom: + fieldRef: + fieldPath: metadata.namespace - name: OSAC_ENABLE_CLUSTER_CONTROLLER value: {{ .Values.controllers.clusterOrder | quote }} - name: OSAC_ENABLE_COMPUTE_INSTANCE_CONTROLLER @@ -106,6 +110,8 @@ spec: value: {{ .Values.controllers.bareMetalInstance | quote }} - name: OSAC_ENABLE_STORAGE_CONTROLLER value: {{ .Values.controllers.storage | quote }} + - name: OSAC_ENABLE_VOLUME_CONTROLLER + value: {{ .Values.controllers.volume | quote }} securityContext: allowPrivilegeEscalation: false readOnlyRootFilesystem: true diff --git a/osac-operator/charts/operator/templates/hub-access-clusterrole.yaml b/osac-operator/charts/operator/templates/hub-access-clusterrole.yaml index e8242f73c3..16841b2e81 100644 --- a/osac-operator/charts/operator/templates/hub-access-clusterrole.yaml +++ b/osac-operator/charts/operator/templates/hub-access-clusterrole.yaml @@ -22,6 +22,7 @@ rules: - subnets - tenants - virtualnetworks + - volumes verbs: - create - delete @@ -44,6 +45,7 @@ rules: - subnets/status - tenants/status - virtualnetworks/status + - volumes/status verbs: - get - apiGroups: diff --git a/osac-operator/charts/operator/values.yaml b/osac-operator/charts/operator/values.yaml index a119e6321a..de14bf87e3 100644 --- a/osac-operator/charts/operator/values.yaml +++ b/osac-operator/charts/operator/values.yaml @@ -31,6 +31,7 @@ controllers: networking: true bareMetalInstance: true storage: true + volume: true configSecret: name: "osac-config" diff --git a/osac-operator/cmd/main.go b/osac-operator/cmd/main.go index 26bd964eac..bf8bab4f3e 100644 --- a/osac-operator/cmd/main.go +++ b/osac-operator/cmd/main.go @@ -82,6 +82,7 @@ const ( envNetworkingNamespace = "OSAC_NETWORKING_NAMESPACE" envClusterOrderNamespace = "OSAC_CLUSTER_ORDER_NAMESPACE" envBareMetalInstanceNamespace = "OSAC_BARE_METAL_INSTANCE_NAMESPACE" + envVolumeNamespace = "OSAC_VOLUME_NAMESPACE" // AAP configuration envAAPURL = "OSAC_AAP_URL" @@ -117,6 +118,7 @@ const ( // Controller enable flags (defaults when flag is not set) envEnableTenantController = "OSAC_ENABLE_TENANT_CONTROLLER" envEnableStorageController = "OSAC_ENABLE_STORAGE_CONTROLLER" + envEnableVolumeController = "OSAC_ENABLE_VOLUME_CONTROLLER" envEnableComputeInstanceController = "OSAC_ENABLE_COMPUTE_INSTANCE_CONTROLLER" envEnableClusterController = "OSAC_ENABLE_CLUSTER_CONTROLLER" envEnableNetworkingController = "OSAC_ENABLE_NETWORKING_CONTROLLER" @@ -135,6 +137,7 @@ const ( type controllerFlags struct { Tenant bool Storage bool + Volume bool ComputeInstance bool Cluster bool Networking bool @@ -150,7 +153,10 @@ func registerControllerFlags() *controllerFlags { "Enable the tenant controller.") flag.BoolVar(&flags.Storage, "enable-storage-controller", helpers.GetEnvWithDefault(envEnableStorageController, false), - "Enable the storage controller.") + "Enable the storage controller (tenant StorageClass management, ClusterOrder storage provisioning).") + flag.BoolVar(&flags.Volume, "enable-volume-controller", + helpers.GetEnvWithDefault(envEnableVolumeController, false), + "Enable the volume controller (block volume provisioning via vendor CSI).") flag.BoolVar(&flags.ComputeInstance, "enable-compute-instance-controller", helpers.GetEnvWithDefault(envEnableComputeInstanceController, false), "Enable the compute-instance controller.") @@ -168,9 +174,10 @@ func registerControllerFlags() *controllerFlags { // enableAllIfNoneSet enables all controllers if none are explicitly enabled. func (f *controllerFlags) enableAllIfNoneSet() { - if !f.Tenant && !f.Storage && !f.ComputeInstance && !f.Cluster && !f.Networking && !f.BareMetalInstance { + if !f.Tenant && !f.Storage && !f.Volume && !f.ComputeInstance && !f.Cluster && !f.Networking && !f.BareMetalInstance { f.Tenant = true f.Storage = true + f.Volume = true f.ComputeInstance = true f.Cluster = true f.Networking = true @@ -471,6 +478,77 @@ func setupStorageController(mgr mcmanager.Manager, grpcConn *grpc.ClientConn, ma return nil } +// setupControllers registers all enabled controllers with the manager. +func setupControllers(mgr mcmanager.Manager, grpcConn *grpc.ClientConn, flags *controllerFlags, maxJobHistory int) error { + if flags.Cluster { + if err := setupClusterControllers(mgr, grpcConn, maxJobHistory); err != nil { + return fmt.Errorf("cluster controllers: %w", err) + } + } + if flags.ComputeInstance { + if err := setupComputeInstanceControllers(mgr, grpcConn, maxJobHistory); err != nil { + return fmt.Errorf("computeinstance controllers: %w", err) + } + } + if flags.Tenant { + if err := setupTenantController(mgr); err != nil { + return fmt.Errorf("tenant controller: %w", err) + } + } + if flags.Storage { + if err := setupStorageController(mgr, grpcConn, maxJobHistory); err != nil { + return fmt.Errorf("storage controller: %w", err) + } + } + if flags.Volume { + if err := setupVolumeControllers(mgr, grpcConn); err != nil { + return fmt.Errorf("volume controllers: %w", err) + } + } + if flags.Networking { + if err := setupNetworkingControllers(mgr, grpcConn, maxJobHistory); err != nil { + return fmt.Errorf("networking controllers: %w", err) + } + } + if flags.BareMetalInstance { + if err := setupBareMetalInstanceControllers(mgr, grpcConn); err != nil { + return fmt.Errorf("baremetalinstance controllers: %w", err) + } + } + return nil +} + +// setupVolumeControllers registers the Volume resource controller and, when +// grpcConn is set, the Volume feedback controller. The Volume controller uses +// a VendorProvisioner interface instead of AAP; for now no real vendor is +// configured (nil provisioner), so the controller sets Progressing and waits +// for the vendor CSI integration in a follow-up PR. +func setupVolumeControllers(mgr mcmanager.Manager, grpcConn *grpc.ClientConn) error { + localMgr := mgr.GetLocalManager() + volumeNamespace := os.Getenv(envVolumeNamespace) + + if grpcConn != nil { + if err := controller.NewVolumeFeedbackReconciler( + localMgr.GetClient(), + grpcConn, + volumeNamespace, + ).SetupWithManager(mgr); err != nil { + return fmt.Errorf("volume feedback controller: %w", err) + } + } + + // VendorProvisioner is nil until the real vendor CSI client is wired. + // The controller will set phase to Progressing and skip provisioning. + if err := controller.NewVolumeReconciler( + mgr, + volumeNamespace, + nil, + ).SetupWithManager(mgr); err != nil { + return fmt.Errorf("volume controller: %w", err) + } + return nil +} + // setupNetworkingControllers registers all networking controllers along with their // feedback controllers when grpcConn is set. func setupNetworkingControllers( @@ -974,41 +1052,9 @@ func main() { }) setupLog.Info("job history configuration", "maxJobs", maxJobHistory) - if ctrlFlags.Cluster { - if err := setupClusterControllers(mgr, grpcConn, maxJobHistory); err != nil { - setupLog.Error(err, "unable to setup cluster controllers") - os.Exit(1) - } - } - if ctrlFlags.ComputeInstance { - if err := setupComputeInstanceControllers(mgr, grpcConn, maxJobHistory); err != nil { - setupLog.Error(err, "unable to setup computeinstance controllers") - os.Exit(1) - } - } - if ctrlFlags.Tenant { - if err := setupTenantController(mgr); err != nil { - setupLog.Error(err, "unable to setup tenant controller") - os.Exit(1) - } - } - if ctrlFlags.Storage { - if err := setupStorageController(mgr, grpcConn, maxJobHistory); err != nil { - setupLog.Error(err, "unable to setup storage controller") - os.Exit(1) - } - } - if ctrlFlags.Networking { - if err := setupNetworkingControllers(mgr, grpcConn, maxJobHistory); err != nil { - setupLog.Error(err, "unable to setup networking controllers") - os.Exit(1) - } - } - if ctrlFlags.BareMetalInstance { - if err := setupBareMetalInstanceControllers(mgr, grpcConn); err != nil { - setupLog.Error(err, "unable to setup baremetalinstance controllers") - os.Exit(1) - } + if err := setupControllers(mgr, grpcConn, ctrlFlags, maxJobHistory); err != nil { + setupLog.Error(err, "unable to setup controllers") + os.Exit(1) } // +kubebuilder:scaffold:builder diff --git a/osac-operator/config/rbac/role.yaml b/osac-operator/config/rbac/role.yaml index 304409c811..18d2581ee0 100644 --- a/osac-operator/config/rbac/role.yaml +++ b/osac-operator/config/rbac/role.yaml @@ -109,6 +109,7 @@ rules: - subnets/finalizers - tenants/finalizers - virtualnetworks/finalizers + - volumes/finalizers verbs: - update - apiGroups: @@ -124,6 +125,7 @@ rules: - subnets - tenants - virtualnetworks + - volumes verbs: - create - delete @@ -145,6 +147,7 @@ rules: - subnets/status - tenants/status - virtualnetworks/status + - volumes/status verbs: - get - patch diff --git a/osac-operator/internal/controller/volume_controller.go b/osac-operator/internal/controller/volume_controller.go new file mode 100644 index 0000000000..d321a19695 --- /dev/null +++ b/osac-operator/internal/controller/volume_controller.go @@ -0,0 +1,298 @@ +/* +Copyright 2026. + +Licensed 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 controller + +import ( + "context" + + "k8s.io/apimachinery/pkg/api/equality" + apimeta "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + controllerutil "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + ctrllog "sigs.k8s.io/controller-runtime/pkg/log" + "sigs.k8s.io/controller-runtime/pkg/predicate" + mcbuilder "sigs.k8s.io/multicluster-runtime/pkg/builder" + mcmanager "sigs.k8s.io/multicluster-runtime/pkg/manager" + mcreconcile "sigs.k8s.io/multicluster-runtime/pkg/reconcile" + + "github.com/osac-project/osac/osac-operator/api/v1alpha1" +) + +// VendorProvisioner abstracts vendor storage array operations. Unlike other +// OSAC resources that provision through AAP (RunProvisioningLifecycle), volumes +// are provisioned by calling the vendor CSI controller directly. This interface +// decouples the controller from the vendor implementation, allowing a mock for +// testing and development while the real vendor CSI client is wired in PR #3. +type VendorProvisioner interface { + CreateVolume(ctx context.Context, req VendorCreateVolumeRequest) (VendorCreateVolumeResponse, error) + DeleteVolume(ctx context.Context, req VendorDeleteVolumeRequest) error +} + +// VendorCreateVolumeRequest carries the parameters the vendor needs to +// provision a volume. Backend is resolved by the fulfillment-service tier +// resolution (OSAC-3277) before the Volume CR is created; the operator +// passes it through to the vendor without re-resolving. +type VendorCreateVolumeRequest struct { + Name string + Backend string + SizeGiB int64 + AccessMode v1alpha1.VolumeAccessMode +} + +// VendorCreateVolumeResponse carries the vendor-assigned identifiers that the +// feedback controller syncs back to the fulfillment-service inventory. +type VendorCreateVolumeResponse struct { + VendorVolumeID string + Backend string + Protocol string +} + +// VendorDeleteVolumeRequest identifies the vendor volume to deprovision. +type VendorDeleteVolumeRequest struct { + VendorVolumeID string + Backend string +} + +// VolumeReconciler reconciles Volume CRs created by the fulfillment-service +// reconciler. It calls the VendorProvisioner to create volumes on the backend +// storage array and updates the CR status with vendor-assigned identifiers. +// The feedback controller then syncs that status back to fulfillment-service. +// +// Unlike networking and compute controllers that use AAP for provisioning, +// this controller calls the vendor CSI directly because storage provisioning +// is a synchronous gRPC call, not an asynchronous job. +type VolumeReconciler struct { + client.Client + Scheme *runtime.Scheme + mgr mcmanager.Manager + VolumeNamespace string + VendorProvisioner VendorProvisioner +} + +// NewVolumeReconciler creates a new reconciler for Volume resources. The +// volumeNamespace controls which namespace the controller watches; in +// production this is set via OSAC_VOLUME_NAMESPACE (same as the Helm release +// namespace), defaulting to "osac-volume" for local development. +func NewVolumeReconciler( + mgr mcmanager.Manager, + volumeNamespace string, + vendorProvisioner VendorProvisioner, +) *VolumeReconciler { + if mgr == nil { + panic("mgr must not be nil") + } + if volumeNamespace == "" { + volumeNamespace = defaultVolumeNamespace + } + return &VolumeReconciler{ + Client: mgr.GetLocalManager().GetClient(), + Scheme: mgr.GetLocalManager().GetScheme(), + mgr: mgr, + VolumeNamespace: volumeNamespace, + VendorProvisioner: vendorProvisioner, + } +} + +// +kubebuilder:rbac:groups=osac.openshift.io,resources=volumes,verbs=get;list;watch;create;update;patch;delete +// +kubebuilder:rbac:groups=osac.openshift.io,resources=volumes/status,verbs=get;update;patch +// +kubebuilder:rbac:groups=osac.openshift.io,resources=volumes/finalizers,verbs=update + +// Reconcile is part of the main Kubernetes reconciliation loop. It drives +// Volume CRs through the provisioning lifecycle: Progressing -> Ready (on +// success) or Failed (on vendor error), and handles deletion by calling +// the vendor to deprovision before removing the finalizer. +func (r *VolumeReconciler) Reconcile(ctx context.Context, req mcreconcile.Request) (ctrl.Result, error) { + log := ctrllog.FromContext(ctx) + + vol := &v1alpha1.Volume{} + if err := r.Get(ctx, req.NamespacedName, vol); err != nil { + return ctrl.Result{}, client.IgnoreNotFound(err) + } + + log.Info("start reconcile") + + oldstatus := vol.Status.DeepCopy() + + var res ctrl.Result + var err error + if vol.ObjectMeta.DeletionTimestamp.IsZero() { + res, err = r.handleUpdate(ctx, vol) + } else { + res, err = r.handleDelete(ctx, vol) + } + + if !equality.Semantic.DeepEqual(vol.Status, *oldstatus) { + log.Info("status requires update") + if err := r.Status().Update(ctx, vol); err != nil { + return res, err + } + } + + log.Info("end reconcile") + return res, err +} + +// handleUpdate runs on every non-deleted reconcile. It ensures the finalizer +// is present, sets the initial phase to Progressing, and delegates to +// handleProvisioning if the volume has not yet reached Ready. +func (r *VolumeReconciler) handleUpdate(ctx context.Context, vol *v1alpha1.Volume) (ctrl.Result, error) { + log := ctrllog.FromContext(ctx) + + if controllerutil.AddFinalizer(vol, osacVolumeFinalizer) { + if err := r.Update(ctx, vol); err != nil { + return ctrl.Result{}, err + } + } + + if vol.Status.Phase == "" { + vol.Status.Phase = v1alpha1.VolumePhaseProgressing + } + + if r.VendorProvisioner == nil { + // Temporary: VendorProvisioner is nil until the real vendor CSI client + // is wired in the CSI driver integration PR. Remove this guard once a + // concrete implementation is always passed to NewVolumeReconciler. + log.Info("no vendor provisioner configured, skipping provisioning") + return ctrl.Result{}, nil + } + + // Already provisioned; nothing to do until spec changes (future: resize). + if vol.Status.Phase == v1alpha1.VolumePhaseReady { + return ctrl.Result{}, nil + } + + // Failed volumes are not auto-retried to avoid spamming the vendor API + // with a persistent configuration error. The admin fixes the issue, then + // a Signal or periodic sync triggers re-reconciliation. At that point the + // phase is reset to Progressing below so provisioning is attempted again. + if vol.Status.Phase == v1alpha1.VolumePhaseFailed { + return ctrl.Result{}, nil + } + + return r.handleProvisioning(ctx, vol) +} + +// handleProvisioning calls the vendor CSI to create the volume on the backend +// array. On success it transitions the phase to Ready and sets the +// VendorProvisioned condition. On failure it transitions to Failed and returns +// nil (no retry) so the error is visible in the condition; the feedback +// controller will sync this state to the fulfillment-service. +func (r *VolumeReconciler) handleProvisioning(ctx context.Context, vol *v1alpha1.Volume) (ctrl.Result, error) { + log := ctrllog.FromContext(ctx) + + resp, err := r.VendorProvisioner.CreateVolume(ctx, VendorCreateVolumeRequest{ + Name: vol.Name, + Backend: vol.Status.Backend, + SizeGiB: vol.Spec.SizeGiB, + AccessMode: vol.Spec.AccessMode, + }) + if err != nil { + log.Error(err, "vendor provisioning failed") + vol.Status.Phase = v1alpha1.VolumePhaseFailed + setVendorProvisionedCondition(&vol.Status.Conditions, metav1.ConditionFalse, "ProvisioningFailed", err.Error()) + return ctrl.Result{}, nil + } + + vol.Status.VendorVolumeID = resp.VendorVolumeID + vol.Status.Backend = resp.Backend + vol.Status.Protocol = v1alpha1.VolumeProtocol(resp.Protocol) + vol.Status.Phase = v1alpha1.VolumePhaseReady + setVendorProvisionedCondition(&vol.Status.Conditions, metav1.ConditionTrue, "Provisioned", "Volume provisioned on vendor storage array") + + log.Info("vendor provisioning succeeded", + "vendorVolumeID", resp.VendorVolumeID, + "backend", resp.Backend, + "protocol", resp.Protocol, + ) + + return ctrl.Result{}, nil +} + +// handleDelete runs when the Volume CR has a deletion timestamp. It calls the +// vendor to deprovision the volume from the backend array, then removes the +// resource controller's finalizer. If deprovisioning fails the error is +// returned so the reconciler retries on the next cycle. +func (r *VolumeReconciler) handleDelete(ctx context.Context, vol *v1alpha1.Volume) (ctrl.Result, error) { + log := ctrllog.FromContext(ctx) + log.Info("deleting volume") + + vol.Status.Phase = v1alpha1.VolumePhaseDeleting + + if !controllerutil.ContainsFinalizer(vol, osacVolumeFinalizer) { + return ctrl.Result{}, nil + } + + // Only call vendor if there is a volume to delete; volumes that failed + // before vendor provisioning have no VendorVolumeID. + if r.VendorProvisioner != nil && vol.Status.VendorVolumeID != "" { + err := r.VendorProvisioner.DeleteVolume(ctx, VendorDeleteVolumeRequest{ + VendorVolumeID: vol.Status.VendorVolumeID, + Backend: vol.Status.Backend, + }) + if err != nil { + log.Error(err, "vendor deprovisioning failed") + return ctrl.Result{}, err + } + log.Info("vendor deprovisioning succeeded", "vendorVolumeID", vol.Status.VendorVolumeID) + } + + if controllerutil.RemoveFinalizer(vol, osacVolumeFinalizer) { + if err := r.Update(ctx, vol); err != nil { + return ctrl.Result{}, err + } + } + + return ctrl.Result{}, nil +} + +// VolumeNamespacePredicate filters events to only those in the configured +// volume namespace, preventing the controller from reacting to Volume CRs +// in other namespaces. +func VolumeNamespacePredicate(namespace string) predicate.Predicate { + return predicate.NewPredicateFuncs( + func(obj client.Object) bool { + return obj.GetNamespace() == namespace + }, + ) +} + +// SetupWithManager registers the Volume controller with the manager. It +// watches Volume CRs in the configured namespace on the local (hub) cluster. +func (r *VolumeReconciler) SetupWithManager(mgr mcmanager.Manager) error { + return mcbuilder.ControllerManagedBy(mgr). + For(&v1alpha1.Volume{}, + mcbuilder.WithPredicates(VolumeNamespacePredicate(r.VolumeNamespace)), + mcbuilder.WithEngageWithLocalCluster(true), + mcbuilder.WithEngageWithProviderClusters(false)). + Complete(r) +} + +// setVendorProvisionedCondition upserts the VendorProvisioned condition. +// Uses apimeta.SetStatusCondition which preserves LastTransitionTime when +// the status hasn't changed, but always updates Reason and Message so +// repeated failures reflect the latest error. +func setVendorProvisionedCondition(conditions *[]metav1.Condition, status metav1.ConditionStatus, reason, message string) { + apimeta.SetStatusCondition(conditions, metav1.Condition{ + Type: string(v1alpha1.VolumeConditionVendorProvisioned), + Status: status, + Reason: reason, + Message: message, + }) +} diff --git a/osac-operator/internal/controller/volume_controller_test.go b/osac-operator/internal/controller/volume_controller_test.go new file mode 100644 index 0000000000..8290e44e1d --- /dev/null +++ b/osac-operator/internal/controller/volume_controller_test.go @@ -0,0 +1,290 @@ +/* +Copyright 2026. + +Licensed 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 controller + +import ( + "context" + "fmt" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + "k8s.io/apimachinery/pkg/api/errors" + apimeta "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + mcreconcile "sigs.k8s.io/multicluster-runtime/pkg/reconcile" + + osacv1alpha1 "github.com/osac-project/osac/osac-operator/api/v1alpha1" +) + +var _ = Describe("VolumeReconciler", func() { + var ( + reconciler *VolumeReconciler + mockProv *MockVendorProvisioner + testCtx context.Context + vol *osacv1alpha1.Volume + ) + + BeforeEach(func() { + testCtx = context.TODO() + mockProv = NewMockVendorProvisioner() + reconciler = &VolumeReconciler{ + Client: k8sClient, + Scheme: k8sClient.Scheme(), + mgr: testMcManager, + VolumeNamespace: "default", + VendorProvisioner: mockProv, + } + + vol = &osacv1alpha1.Volume{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-vol", + Namespace: "default", + }, + Spec: osacv1alpha1.VolumeSpec{ + StorageTier: "gold", + SizeGiB: 100, + AccessMode: osacv1alpha1.VolumeAccessModeReadWriteOnce, + }, + } + }) + + AfterEach(func() { + volKey := types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace} + existingVol := &osacv1alpha1.Volume{} + if err := k8sClient.Get(testCtx, volKey, existingVol); err == nil { + existingVol.Finalizers = nil + _ = k8sClient.Update(testCtx, existingVol) + _ = k8sClient.Delete(testCtx, existingVol) + } + }) + + It("should add finalizer on first reconcile", func() { + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + + updated := &osacv1alpha1.Volume{} + Expect(k8sClient.Get(testCtx, types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, updated)).To(Succeed()) + Expect(updated.Finalizers).To(ContainElement(osacVolumeFinalizer)) + }) + + It("should reach Ready on first reconcile when the mock provisioner succeeds", func() { + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + + updated := &osacv1alpha1.Volume{} + Expect(k8sClient.Get(testCtx, types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, updated)).To(Succeed()) + // Phase should be Ready because the mock provisioner succeeds immediately + Expect(updated.Status.Phase).To(Equal(osacv1alpha1.VolumePhaseReady)) + }) + + It("should provision volume and set status fields on success", func() { + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + // First reconcile adds finalizer + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + + // Second reconcile provisions + _, err = reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + + updated := &osacv1alpha1.Volume{} + Expect(k8sClient.Get(testCtx, types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, updated)).To(Succeed()) + + Expect(updated.Status.Phase).To(Equal(osacv1alpha1.VolumePhaseReady)) + Expect(updated.Status.VendorVolumeID).To(HavePrefix("mock-")) + Expect(updated.Status.Backend).To(Equal("mock-backend")) + Expect(updated.Status.Protocol).To(Equal(osacv1alpha1.VolumeProtocolBlock)) + Expect(mockProv.CreateCallCount()).To(BeNumerically(">=", 1)) + + cond := apimeta.FindStatusCondition(updated.Status.Conditions, string(osacv1alpha1.VolumeConditionVendorProvisioned)) + Expect(cond).ToNot(BeNil()) + Expect(cond.Status).To(Equal(metav1.ConditionTrue)) + Expect(cond.Reason).To(Equal("Provisioned")) + }) + + It("should set phase to Failed when vendor provisioning fails", func() { + mockProv.CreateErr = fmt.Errorf("vendor array unreachable") + + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + // First reconcile adds finalizer + _, _ = reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + + // Second reconcile attempts provisioning + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + + updated := &osacv1alpha1.Volume{} + Expect(k8sClient.Get(testCtx, types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, updated)).To(Succeed()) + + Expect(updated.Status.Phase).To(Equal(osacv1alpha1.VolumePhaseFailed)) + + cond := apimeta.FindStatusCondition(updated.Status.Conditions, string(osacv1alpha1.VolumeConditionVendorProvisioned)) + Expect(cond).ToNot(BeNil()) + Expect(cond.Status).To(Equal(metav1.ConditionFalse)) + Expect(cond.Reason).To(Equal("ProvisioningFailed")) + Expect(cond.Message).To(ContainSubstring("vendor array unreachable")) + }) + + It("should not re-provision when already Ready", func() { + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + // Reconcile until Ready + for range 3 { + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + } + + countBefore := mockProv.CreateCallCount() + + // One more reconcile should be a no-op + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + Expect(mockProv.CreateCallCount()).To(Equal(countBefore)) + }) + + It("should handle deletion with vendor deprovisioning", func() { + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + // Reconcile to Ready + for range 3 { + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + } + + // Delete + Expect(k8sClient.Delete(testCtx, vol)).To(Succeed()) + + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + Expect(mockProv.DeleteCallCount()).To(BeNumerically(">=", 1)) + + // Volume should be gone after finalizer removal + deleted := &osacv1alpha1.Volume{} + err = k8sClient.Get(testCtx, types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, deleted) + Expect(errors.IsNotFound(err)).To(BeTrue()) + }) + + It("should return error and keep finalizer when vendor deprovisioning fails", func() { + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + // Reconcile to Ready + for range 3 { + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + } + + // Inject vendor delete failure + mockProv.DeleteErr = fmt.Errorf("storage array unavailable") + + Expect(k8sClient.Delete(testCtx, vol)).To(Succeed()) + + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + Expect(err).To(HaveOccurred()) + + // Finalizer must still be present — the volume was not deprovisioned + still := &osacv1alpha1.Volume{} + Expect(k8sClient.Get(testCtx, types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, still)).To(Succeed()) + Expect(still.Finalizers).To(ContainElement(osacVolumeFinalizer)) + // Phase must be Deleting so the feedback controller syncs VOLUME_STATE_DELETING + Expect(still.Status.Phase).To(Equal(osacv1alpha1.VolumePhaseDeleting)) + }) + + It("should return not-found gracefully when volume is already deleted", func() { + _, err := reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: "nonexistent", Namespace: "default"}, + }, + }) + Expect(err).ToNot(HaveOccurred()) + }) + + It("should skip provisioning when no VendorProvisioner is configured", func() { + reconciler.VendorProvisioner = nil + + Expect(k8sClient.Create(testCtx, vol)).To(Succeed()) + + // Reconcile adds finalizer + sets Progressing + for range 2 { + _, _ = reconciler.Reconcile(testCtx, mcreconcile.Request{ + Request: reconcile.Request{ + NamespacedName: types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, + }, + }) + } + + updated := &osacv1alpha1.Volume{} + Expect(k8sClient.Get(testCtx, types.NamespacedName{Name: vol.Name, Namespace: vol.Namespace}, updated)).To(Succeed()) + Expect(updated.Status.Phase).To(Equal(osacv1alpha1.VolumePhaseProgressing)) + Expect(updated.Status.VendorVolumeID).To(BeEmpty()) + }) +}) diff --git a/osac-operator/internal/controller/volume_feedback_controller.go b/osac-operator/internal/controller/volume_feedback_controller.go new file mode 100644 index 0000000000..7fc689be95 --- /dev/null +++ b/osac-operator/internal/controller/volume_feedback_controller.go @@ -0,0 +1,183 @@ +/* +Copyright (c) 2026 Red Hat Inc. + +Licensed 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 controller + +import ( + "context" + "errors" + "fmt" + + "google.golang.org/grpc" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" + clnt "sigs.k8s.io/controller-runtime/pkg/client" + ctrllog "sigs.k8s.io/controller-runtime/pkg/log" + mcmanager "sigs.k8s.io/multicluster-runtime/pkg/manager" + + "github.com/osac-project/osac/osac-operator/api/v1alpha1" + privatev1 "github.com/osac-project/osac/osac-operator/internal/api/osac/private/v1" + "github.com/osac-project/osac/osac-operator/internal/controller/feedback" +) + +// VolumeFeedbackReconciler syncs Volume CR status from the hub cluster back +// to the fulfillment-service via the private Volumes gRPC API. It maps CRD +// phases to proto states and copies vendor-assigned fields (vendorVolumeID, +// backend, protocol) so the fulfillment-service inventory stays current. +type VolumeFeedbackReconciler struct { + bridge *feedback.Bridge[*v1alpha1.Volume, *privatev1.Volume] + volumeNamespace string +} + +// NewVolumeFeedbackReconciler creates a feedback reconciler that syncs Volume +// CR status to the fulfillment-service. The volumeNamespace controls which +// namespace this controller watches, matching the resource controller's scope. +func NewVolumeFeedbackReconciler(hubClient clnt.Client, grpcConn *grpc.ClientConn, volumeNamespace string) *VolumeFeedbackReconciler { + if volumeNamespace == "" { + volumeNamespace = defaultVolumeNamespace + } + volClient := privatev1.NewVolumesClient(grpcConn) + r := &VolumeFeedbackReconciler{volumeNamespace: volumeNamespace} + r.bridge = &feedback.Bridge[*v1alpha1.Volume, *privatev1.Volume]{ + Client: hubClient, + Finalizer: osacVolumeFeedbackFinalizer, + IDLabel: osacVolumeIDLabel, + Kind: "Volume", + IDKey: "volumeID", + NewObject: func() *v1alpha1.Volume { return &v1alpha1.Volume{} }, + Fetch: func(ctx context.Context, id string) (*privatev1.Volume, error) { + response, err := volClient.Get(ctx, privatev1.VolumesGetRequest_builder{Id: id}.Build()) + if err != nil { + return nil, err + } + vol := response.GetObject() + if vol == nil { + return nil, errors.New("volume response contained nil object") + } + if !vol.HasSpec() { + vol.SetSpec(&privatev1.VolumeSpec{}) + } + if !vol.HasStatus() { + vol.SetStatus(&privatev1.VolumeStatus{}) + } + return vol, nil + }, + Save: func(ctx context.Context, remote *privatev1.Volume) error { + _, err := volClient.Update(ctx, privatev1.VolumesUpdateRequest_builder{ + Object: remote, + }.Build()) + return err + }, + Signal: func(ctx context.Context, id string) error { + _, err := volClient.Signal(ctx, privatev1.VolumesSignalRequest_builder{ + Id: id, + }.Build()) + return err + }, + SyncUpdate: syncVolumeUpdate, + SyncDelete: syncVolumeDelete, + } + return r +} + +// SetupWithManager registers the feedback controller with the manager. It +// watches Volume CRs in the configured namespace on the local (hub) cluster. +func (r *VolumeFeedbackReconciler) SetupWithManager(mgr mcmanager.Manager) error { + localMgr := mgr.GetLocalManager() + if localMgr == nil { + return fmt.Errorf("local manager is nil") + } + + return ctrl.NewControllerManagedBy(localMgr). + Named("volume-feedback"). + For(&v1alpha1.Volume{}, builder.WithPredicates(VolumeNamespacePredicate(r.volumeNamespace))). + Complete(r) +} + +// Reconcile delegates to the shared feedback Bridge which handles the +// finalizer lifecycle, clone-compare-save, and last-finalizer Signal. +func (r *VolumeFeedbackReconciler) Reconcile(ctx context.Context, request ctrl.Request) (ctrl.Result, error) { + return r.bridge.Reconcile(ctx, request) +} + +// syncVolumeUpdate maps Volume CR status to the fulfillment-service proto on +// the non-delete path. It syncs the phase, vendor-assigned identifiers, and +// the PVC/PV references that the operator populates after provisioning. +func syncVolumeUpdate(ctx context.Context, obj *v1alpha1.Volume, remote *privatev1.Volume) error { + syncVolumePhase(ctx, obj, remote) + syncVolumeVendorFields(obj, remote) + return nil +} + +// syncVolumeDelete maps Volume CR status during deletion. Failed volumes +// report FAILED; all other deletion states report DELETING. +func syncVolumeDelete(_ context.Context, obj *v1alpha1.Volume, remote *privatev1.Volume) error { + if obj.Status.Phase == v1alpha1.VolumePhaseFailed { + remote.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_FAILED) + return nil + } + remote.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_DELETING) + return nil +} + +// syncVolumePhase converts the CRD phase to the proto state enum. +// +// CRD Phase -> Proto State +// Progressing -> CREATING (volume is being provisioned on vendor array) +// Ready -> AVAILABLE (vendor provisioned, ready for use) +// Failed -> FAILED (vendor provisioning failed) +// Deleting -> DELETING (volume is being deprovisioned) +func syncVolumePhase(ctx context.Context, obj *v1alpha1.Volume, remote *privatev1.Volume) { + switch obj.Status.Phase { + case v1alpha1.VolumePhaseProgressing: + remote.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_CREATING) + case v1alpha1.VolumePhaseReady: + remote.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_AVAILABLE) + case v1alpha1.VolumePhaseFailed: + remote.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_FAILED) + case v1alpha1.VolumePhaseDeleting: + remote.GetStatus().SetState(privatev1.VolumeState_VOLUME_STATE_DELETING) + default: + log := ctrllog.FromContext(ctx) + log.Info("Unknown phase, will ignore it", "phase", obj.Status.Phase) + } +} + +// syncVolumeVendorFields copies the vendor-assigned identifiers from the CR +// status to the proto status so the fulfillment-service inventory reflects +// the actual storage array state. +func syncVolumeVendorFields(obj *v1alpha1.Volume, remote *privatev1.Volume) { + if obj.Status.VendorVolumeID != "" { + remote.GetStatus().SetVendorVolumeId(obj.Status.VendorVolumeID) + } + if obj.Status.Backend != "" { + remote.GetStatus().SetBackend(obj.Status.Backend) + } + if obj.Status.Protocol != "" { + remote.GetStatus().SetProtocol(crdProtocolToProto(obj.Status.Protocol)) + } +} + +// crdProtocolToProto converts the CRD VolumeProtocol typed string (e.g. +// "Block", "NFS") to the proto StorageProtocol enum. A direct map lookup +// would fail because the proto keys are "STORAGE_PROTOCOL_BLOCK", not "Block". +func crdProtocolToProto(protocol v1alpha1.VolumeProtocol) privatev1.StorageProtocol { + switch protocol { + case v1alpha1.VolumeProtocolBlock: + return privatev1.StorageProtocol_STORAGE_PROTOCOL_BLOCK + case v1alpha1.VolumeProtocolNFS: + return privatev1.StorageProtocol_STORAGE_PROTOCOL_NFS + default: + return privatev1.StorageProtocol_STORAGE_PROTOCOL_UNSPECIFIED + } +} diff --git a/osac-operator/internal/controller/volume_feedback_controller_test.go b/osac-operator/internal/controller/volume_feedback_controller_test.go new file mode 100644 index 0000000000..580da1577f --- /dev/null +++ b/osac-operator/internal/controller/volume_feedback_controller_test.go @@ -0,0 +1,484 @@ +/* +Copyright 2026. + +Licensed 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 controller + +import ( + "context" + "fmt" + "net" + "sync" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/credentials/insecure" + grpcstatus "google.golang.org/grpc/status" + "google.golang.org/grpc/test/bufconn" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + "github.com/osac-project/osac/osac-operator/api/v1alpha1" + privatev1 "github.com/osac-project/osac/osac-operator/internal/api/osac/private/v1" +) + +var _ = Describe("VolumeFeedbackController", func() { + const ( + volName = "test-vol" + volNamespace = "test-namespace" + volID = "vol-123" + ) + + var ( + ctx context.Context + fakeK8s client.Client + mockServer *mockVolumesServer + reconciler *VolumeFeedbackReconciler + grpcServer *grpc.Server + listener *bufconn.Listener + grpcConn *grpc.ClientConn + ) + + BeforeEach(func() { + ctx = context.Background() + + scheme := runtime.NewScheme() + Expect(v1alpha1.AddToScheme(scheme)).To(Succeed()) + fakeK8s = fake.NewClientBuilder().WithScheme(scheme).Build() + + mockServer = &mockVolumesServer{ + volumes: make(map[string]*privatev1.Volume), + updates: make([]*privatev1.Volume, 0), + signals: make([]string, 0), + } + listener = bufconn.Listen(1024 * 1024) + grpcServer = grpc.NewServer() + privatev1.RegisterVolumesServer(grpcServer, mockServer) + + go func() { + _ = grpcServer.Serve(listener) + }() + + var err error + grpcConn, err = grpc.NewClient("passthrough:///bufnet", + grpc.WithContextDialer(func(ctx context.Context, s string) (net.Conn, error) { + return listener.Dial() + }), + grpc.WithTransportCredentials(insecure.NewCredentials()), + ) + Expect(err).NotTo(HaveOccurred()) + + reconciler = NewVolumeFeedbackReconciler(fakeK8s, grpcConn, volNamespace) + }) + + AfterEach(func() { + if grpcConn != nil { + _ = grpcConn.Close() + } + if grpcServer != nil { + grpcServer.Stop() + } + if listener != nil { + _ = listener.Close() + } + }) + + Context("phase-to-state mapping", func() { + It("should sync Phase=Ready to state=AVAILABLE", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_CREATING)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseReady, nil) + cr.Status.VendorVolumeID = "vast-001" + cr.Status.Backend = "vast-backend" + cr.Status.Protocol = v1alpha1.VolumeProtocolBlock + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + updated := mockServer.updates[0] + Expect(updated.GetStatus().GetState()).To(Equal(privatev1.VolumeState_VOLUME_STATE_AVAILABLE)) + Expect(updated.GetStatus().GetVendorVolumeId()).To(Equal("vast-001")) + Expect(updated.GetStatus().GetBackend()).To(Equal("vast-backend")) + + // Signal should not be called on non-delete reconciles + Expect(mockServer.signals).To(BeEmpty()) + + updatedCR := &v1alpha1.Volume{} + Expect(fakeK8s.Get(ctx, types.NamespacedName{Name: volName, Namespace: volNamespace}, updatedCR)).To(Succeed()) + Expect(controllerutil.ContainsFinalizer(updatedCR, osacVolumeFeedbackFinalizer)).To(BeTrue()) + }) + + It("should sync Phase=Progressing to state=CREATING", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_AVAILABLE)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseProgressing, nil) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + Expect(mockServer.updates[0].GetStatus().GetState()).To(Equal(privatev1.VolumeState_VOLUME_STATE_CREATING)) + Expect(mockServer.signals).To(BeEmpty()) + }) + + It("should sync Phase=Failed to state=FAILED", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_CREATING)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseFailed, nil) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + Expect(mockServer.updates[0].GetStatus().GetState()).To(Equal(privatev1.VolumeState_VOLUME_STATE_FAILED)) + Expect(mockServer.signals).To(BeEmpty()) + }) + }) + + Context("vendor field syncing", func() { + It("should sync vendorVolumeID, backend, and protocol to remote", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_CREATING)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseReady, nil) + cr.Status.VendorVolumeID = "netapp-vol-42" + cr.Status.Backend = "netapp-cluster-1" + cr.Status.Protocol = v1alpha1.VolumeProtocolNFS + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + updated := mockServer.updates[0] + Expect(updated.GetStatus().GetVendorVolumeId()).To(Equal("netapp-vol-42")) + Expect(updated.GetStatus().GetBackend()).To(Equal("netapp-cluster-1")) + Expect(updated.GetStatus().GetProtocol()).To(Equal(privatev1.StorageProtocol_STORAGE_PROTOCOL_NFS)) + }) + + It("should map Block protocol correctly", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_CREATING)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseReady, nil) + cr.Status.VendorVolumeID = "vast-001" + cr.Status.Backend = "vast-backend" + cr.Status.Protocol = v1alpha1.VolumeProtocolBlock + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + Expect(mockServer.updates[0].GetStatus().GetProtocol()).To(Equal(privatev1.StorageProtocol_STORAGE_PROTOCOL_BLOCK)) + }) + + It("should not overwrite remote fields when CR fields are empty", func() { + remote := newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_CREATING) + remote.GetStatus().SetVendorVolumeId("existing-id") + remote.GetStatus().SetBackend("existing-backend") + mockServer.addVolume(remote) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseProgressing, nil) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + // Progressing maps to CREATING, which the remote already reports, and the CR + // carries no vendor fields. Nothing changes, so no Update RPC is sent. + Expect(mockServer.updates).To(BeEmpty()) + }) + }) + + Context("label and identity handling", func() { + It("should skip CRs without volume-uuid label", func() { + cr := &v1alpha1.Volume{ + ObjectMeta: metav1.ObjectMeta{ + Name: volName, + Namespace: volNamespace, + Labels: map[string]string{}, + }, + Spec: v1alpha1.VolumeSpec{ + StorageTier: "gold", + SizeGiB: 100, + AccessMode: v1alpha1.VolumeAccessModeReadWriteOnce, + }, + Status: v1alpha1.VolumeStatus{ + Phase: v1alpha1.VolumePhaseReady, + }, + } + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + Expect(mockServer.updates).To(BeEmpty()) + }) + }) + + Context("deletion handling", func() { + It("should sync Phase=Deleting to state=DELETING during deletion", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_AVAILABLE)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseDeleting, + []string{osacVolumeFeedbackFinalizer, osacVolumeFinalizer}) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + Expect(fakeK8s.Delete(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + Expect(mockServer.updates[0].GetStatus().GetState()).To(Equal(privatev1.VolumeState_VOLUME_STATE_DELETING)) + + // Signal should NOT be called when other finalizers remain + Expect(mockServer.signals).To(BeEmpty()) + + // Feedback finalizer should remain (other finalizers still present) + updatedCR := &v1alpha1.Volume{} + Expect(fakeK8s.Get(ctx, types.NamespacedName{Name: volName, Namespace: volNamespace}, updatedCR)).To(Succeed()) + Expect(controllerutil.ContainsFinalizer(updatedCR, osacVolumeFeedbackFinalizer)).To(BeTrue()) + }) + + It("should sync Phase=Failed to state=FAILED during deletion", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_AVAILABLE)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseFailed, + []string{osacVolumeFeedbackFinalizer, osacVolumeFinalizer}) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + Expect(fakeK8s.Delete(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + Expect(mockServer.updates[0].GetStatus().GetState()).To(Equal(privatev1.VolumeState_VOLUME_STATE_FAILED)) + }) + + It("should remove finalizer and signal when feedback finalizer is the last one", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_AVAILABLE)) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseDeleting, + []string{osacVolumeFeedbackFinalizer}) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + Expect(fakeK8s.Delete(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(HaveLen(1)) + Expect(mockServer.updates[0].GetStatus().GetState()).To(Equal(privatev1.VolumeState_VOLUME_STATE_DELETING)) + Expect(mockServer.signals).To(HaveLen(1)) + Expect(mockServer.signals[0]).To(Equal(volID)) + + // CR should be gone (last finalizer removed) + updatedCR := &v1alpha1.Volume{} + err = fakeK8s.Get(ctx, types.NamespacedName{Name: volName, Namespace: volNamespace}, updatedCR) + Expect(err).To(HaveOccurred()) + }) + + It("should still remove finalizer when signal fails", func() { + mockServer.addVolume(newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_AVAILABLE)) + mockServer.signalErr = fmt.Errorf("signal unavailable") + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseDeleting, + []string{osacVolumeFeedbackFinalizer}) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + Expect(fakeK8s.Delete(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + // CR should still be gone (finalizer removed despite signal failure) + updatedCR := &v1alpha1.Volume{} + err = fakeK8s.Get(ctx, types.NamespacedName{Name: volName, Namespace: volNamespace}, updatedCR) + Expect(err).To(HaveOccurred()) + }) + + It("should remove feedback finalizer when remote record is NotFound during deletion", func() { + // Don't add volume to mock server (simulates archived record) + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseDeleting, + []string{osacVolumeFeedbackFinalizer}) + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + Expect(fakeK8s.Delete(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + Expect(mockServer.updates).To(BeEmpty()) + Expect(mockServer.signals).To(BeEmpty()) + + // CR should be gone + updatedCR := &v1alpha1.Volume{} + err = fakeK8s.Get(ctx, types.NamespacedName{Name: volName, Namespace: volNamespace}, updatedCR) + Expect(err).To(HaveOccurred()) + }) + }) + + Context("idempotency", func() { + It("should not call Update when remote state already matches", func() { + remote := newRemoteVolume(volID, privatev1.VolumeState_VOLUME_STATE_AVAILABLE) + remote.GetStatus().SetVendorVolumeId("vast-001") + remote.GetStatus().SetBackend("vast-backend") + mockServer.addVolume(remote) + + cr := newVolumeFeedbackCR(volName, volNamespace, volID, v1alpha1.VolumePhaseReady, nil) + cr.Status.VendorVolumeID = "vast-001" + cr.Status.Backend = "vast-backend" + // Pre-seed the feedback finalizer so the reconciler doesn't add it (which triggers an update) + cr.Finalizers = []string{osacVolumeFeedbackFinalizer} + Expect(fakeK8s.Create(ctx, cr)).To(Succeed()) + + _, err := reconciler.Reconcile(ctx, reconcile.Request{ + NamespacedName: types.NamespacedName{Name: volName, Namespace: volNamespace}, + }) + Expect(err).NotTo(HaveOccurred()) + + // No Update RPC called since remote already matches + Expect(mockServer.updates).To(BeEmpty()) + }) + }) +}) + +// Helper to create a Volume CR for feedback controller tests. +func newVolumeFeedbackCR(name, namespace, id string, phase v1alpha1.VolumePhaseType, finalizers []string) *v1alpha1.Volume { + cr := &v1alpha1.Volume{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: namespace, + Labels: map[string]string{ + osacVolumeIDLabel: id, + }, + }, + Spec: v1alpha1.VolumeSpec{ + StorageTier: "gold", + SizeGiB: 100, + AccessMode: v1alpha1.VolumeAccessModeReadWriteOnce, + }, + Status: v1alpha1.VolumeStatus{ + Phase: phase, + }, + } + if len(finalizers) > 0 { + cr.Finalizers = finalizers + } + return cr +} + +// Helper to create a remote Volume proto for the mock server. +func newRemoteVolume(id string, state privatev1.VolumeState) *privatev1.Volume { + return privatev1.Volume_builder{ + Id: id, + Metadata: privatev1.Metadata_builder{ + Name: "test-vol", + }.Build(), + Spec: privatev1.VolumeSpec_builder{ + StorageTier: "gold", + SizeGib: 100, + AccessMode: privatev1.VolumeAccessMode_VOLUME_ACCESS_MODE_READ_WRITE_ONCE, + }.Build(), + Status: privatev1.VolumeStatus_builder{ + State: state, + }.Build(), + }.Build() +} + +// mockVolumesServer implements privatev1.VolumesServer for testing. +type mockVolumesServer struct { + privatev1.UnimplementedVolumesServer + mu sync.Mutex + volumes map[string]*privatev1.Volume + updates []*privatev1.Volume + signals []string + signalErr error +} + +func (m *mockVolumesServer) addVolume(vol *privatev1.Volume) { + m.mu.Lock() + defer m.mu.Unlock() + m.volumes[vol.GetId()] = vol +} + +func (m *mockVolumesServer) Get(_ context.Context, req *privatev1.VolumesGetRequest) (*privatev1.VolumesGetResponse, error) { + m.mu.Lock() + defer m.mu.Unlock() + + vol, ok := m.volumes[req.GetId()] + if !ok { + return nil, grpcstatus.Errorf(codes.NotFound, "object with identifier '%s' not found", req.GetId()) + } + + return privatev1.VolumesGetResponse_builder{ + Object: vol, + }.Build(), nil +} + +func (m *mockVolumesServer) Update(_ context.Context, req *privatev1.VolumesUpdateRequest) (*privatev1.VolumesUpdateResponse, error) { + m.mu.Lock() + defer m.mu.Unlock() + + vol := req.GetObject() + m.volumes[vol.GetId()] = vol + m.updates = append(m.updates, vol) + + return privatev1.VolumesUpdateResponse_builder{ + Object: vol, + }.Build(), nil +} + +func (m *mockVolumesServer) Signal(_ context.Context, req *privatev1.VolumesSignalRequest) (*privatev1.VolumesSignalResponse, error) { + m.mu.Lock() + defer m.mu.Unlock() + + m.signals = append(m.signals, req.GetId()) + + if m.signalErr != nil { + return nil, m.signalErr + } + return &privatev1.VolumesSignalResponse{}, nil +} diff --git a/osac-operator/internal/controller/volume_mock_provisioner_test.go b/osac-operator/internal/controller/volume_mock_provisioner_test.go new file mode 100644 index 0000000000..01f6745cc3 --- /dev/null +++ b/osac-operator/internal/controller/volume_mock_provisioner_test.go @@ -0,0 +1,75 @@ +/* +Copyright 2026. + +Licensed 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 controller + +import ( + "context" + "fmt" + "sync/atomic" +) + +// MockVendorProvisioner is a test-only VendorProvisioner that succeeds +// immediately with deterministic IDs. It tracks call counts for test +// assertions. Set CreateErr or DeleteErr to simulate vendor failures. +type MockVendorProvisioner struct { + createCount atomic.Int64 + deleteCount atomic.Int64 + + // CreateErr, when non-nil, is returned by CreateVolume instead of + // succeeding. Allows tests to simulate vendor failures. + CreateErr error + + // DeleteErr, when non-nil, is returned by DeleteVolume instead of + // succeeding. + DeleteErr error +} + +// NewMockVendorProvisioner creates a mock provisioner that succeeds by default. +func NewMockVendorProvisioner() *MockVendorProvisioner { + return &MockVendorProvisioner{} +} + +// CreateVolume returns a deterministic vendor volume ID composed of +// "mock-" plus a monotonic counter. Backend and protocol are fixed +// strings suitable for test assertions. +func (m *MockVendorProvisioner) CreateVolume(_ context.Context, req VendorCreateVolumeRequest) (VendorCreateVolumeResponse, error) { + n := m.createCount.Add(1) + if m.CreateErr != nil { + return VendorCreateVolumeResponse{}, m.CreateErr + } + return VendorCreateVolumeResponse{ + VendorVolumeID: fmt.Sprintf("mock-%d", n), + Backend: "mock-backend", + Protocol: "Block", + }, nil +} + +// DeleteVolume records the call and returns DeleteErr (nil by default). +func (m *MockVendorProvisioner) DeleteVolume(_ context.Context, _ VendorDeleteVolumeRequest) error { + m.deleteCount.Add(1) + return m.DeleteErr +} + +// CreateCallCount returns the number of times CreateVolume was called. +func (m *MockVendorProvisioner) CreateCallCount() int64 { + return m.createCount.Load() +} + +// DeleteCallCount returns the number of times DeleteVolume was called. +func (m *MockVendorProvisioner) DeleteCallCount() int64 { + return m.deleteCount.Load() +} diff --git a/osac-operator/internal/controller/volume_names.go b/osac-operator/internal/controller/volume_names.go new file mode 100644 index 0000000000..f581c9ea4c --- /dev/null +++ b/osac-operator/internal/controller/volume_names.go @@ -0,0 +1,31 @@ +/* +Copyright 2026. + +Licensed 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 controller + +import ( + "fmt" +) + +const ( + defaultVolumeNamespace = "osac-volume" +) + +var ( + osacVolumeIDLabel string = fmt.Sprintf("%s/volume-uuid", osacPrefix) + osacVolumeFinalizer string = fmt.Sprintf("%s/volume-finalizer", osacPrefix) + osacVolumeFeedbackFinalizer string = fmt.Sprintf("%s/volume-feedback", osacPrefix) +)