Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions pkg/driver/gcp-pd/gcp_pd.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,14 @@ func GetGCPPDOperatorConfig() *config.OperatorConfig {
AssetDir: generatedAssetBase,
OperatorControllerConfigBuilder: GetGCPPDOperatorControllerConfig,
Removable: false,
PrerequisiteAssets: []string{
"controller_sa.yaml",
"node_sa.yaml",
"hostnetwork_role.yaml",
"controller_hostnetwork_binding.yaml",
"privileged_role.yaml",
"node_privileged_binding.yaml",
},
}
}

Expand Down
2 changes: 2 additions & 0 deletions pkg/operator/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ type OperatorConfig struct {
CloudConfigNamespace string
// Removable should be true if the operator and its operand can be removed
Removable bool
// Prerequisite (static) assets
PrerequisiteAssets []string
}

// OperatorControllerConfig is configuration of controllers that are used to deploy CSI drivers.
Expand Down
71 changes: 71 additions & 0 deletions pkg/operator/prerequisites.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
package operator

import (
"context"
"fmt"
"path/filepath"
"time"

"github.com/openshift/csi-operator/assets"
"github.com/openshift/library-go/pkg/operator/events"
"github.com/openshift/library-go/pkg/operator/resource/resourceapply"
kubeclient "k8s.io/client-go/kubernetes"
"k8s.io/klog/v2"
)

var (
numIterations = 10
delayIteration = 5 * time.Second
)

func applyPrerequisites(ctx context.Context, kubeClient kubeclient.Interface, recorder events.Recorder, assetDir string, assetNames []string) error {
files := make([]string, len(assetNames))
for i, name := range assetNames {
files[i] = filepath.Join(assetDir, name)

}

var errs []error
for range numIterations {
if len(files) == 0 {
klog.Infof("All prerequisite assets are applied")
return nil
}
results := resourceapply.ApplyDirectly(
ctx,
resourceapply.NewKubeClientHolder(kubeClient),
recorder,
resourceapply.NewResourceCache(),
assets.ReadFile,
files...,
)

errs = errs[:0]
for _, result := range results {
if result.Error != nil {
errs = append(errs, fmt.Errorf("%s: %w", result.File, result.Error))
continue
}
klog.V(2).Infof("Applied prerequisite asset %s (changed=%v)", result.File, result.Changed)
// Remove successfully applied assets from the list
for i, file := range files {
if file == result.File {
files = append(files[:i], files[i+1:]...)
break
}
}
}
if len(errs) != 0 {
klog.Warningf("Failed to apply some prerequisites: %v", errs)
}
if len(files) != 0 {
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(delayIteration):
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
return fmt.Errorf("failed to apply some prerequisites: %v", errs)

}
150 changes: 150 additions & 0 deletions pkg/operator/prerequisites_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
package operator

import (
"context"
"fmt"
"strings"
"testing"

"github.com/openshift/library-go/pkg/operator/events"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
fakecore "k8s.io/client-go/kubernetes/fake"
k8stesting "k8s.io/client-go/testing"
"k8s.io/utils/clock"
)

func TestApplyPrerequisites(t *testing.T) {
assetDir := "overlays/gcp-pd/generated/standalone"

cases := []struct {
name string
assetDir string
assetNames []string
setupReactor func(client *fakecore.Clientset)
expectErr bool
errContains string
}{
{
name: "should apply all prerequisite assets successfully",
assetDir: assetDir,
assetNames: []string{"controller_sa.yaml", "node_sa.yaml"},
expectErr: false,
},
{
name: "should succeed with empty asset list",
assetDir: assetDir,
assetNames: []string{},
expectErr: false,
},
{
name: "should apply all GCP PD prerequisite assets",
assetDir: assetDir,
assetNames: []string{
"controller_sa.yaml",
"node_sa.yaml",
"hostnetwork_role.yaml",
"controller_hostnetwork_binding.yaml",
"privileged_role.yaml",
"node_privileged_binding.yaml",
},
expectErr: false,
},
{
name: "should return error for non-existent asset file",
assetDir: assetDir,
assetNames: []string{"nonexistent.yaml"},
expectErr: true,
errContains: "nonexistent.yaml",
},
{
name: "should return error when some assets do not exist",
assetDir: assetDir,
assetNames: []string{"controller_sa.yaml", "nonexistent.yaml"},
expectErr: true,
errContains: "nonexistent.yaml",
},
{
name: "should retry failed assets and eventually succeed",
assetDir: assetDir,
assetNames: []string{
"controller_sa.yaml",
"node_sa.yaml",
},
setupReactor: func(client *fakecore.Clientset) {
callCount := 0
client.PrependReactor("create", "serviceaccounts", func(action k8stesting.Action) (bool, runtime.Object, error) {
createAction := action.(k8stesting.CreateAction)
sa := createAction.GetObject().(*corev1.ServiceAccount)
if sa.Name == "gcp-pd-csi-driver-node-sa" {
callCount++
if callCount <= 1 {
return true, nil, fmt.Errorf("transient error")
}
}
return false, nil, nil
})
},
expectErr: false,
},
{
name: "should return error for non-existent asset directory",
assetDir: "no/such/dir",
assetNames: []string{"controller_sa.yaml"},
expectErr: true,
},
}

for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
origDelay := delayIteration
origIterations := numIterations
delayIteration = 0
numIterations = 3
defer func() {
delayIteration = origDelay
numIterations = origIterations
}()

kubeClient := fakecore.NewClientset()
if tc.setupReactor != nil {
tc.setupReactor(kubeClient)
}
recorder := events.NewInMemoryRecorder("test", &clock.RealClock{})
ctx := context.Background()

err := applyPrerequisites(ctx, kubeClient, recorder, tc.assetDir, tc.assetNames)
if tc.expectErr && err == nil {
t.Fatalf("expected error but got nil")
}
if !tc.expectErr && err != nil {
t.Fatalf("unexpected error: %v", err)
}
if tc.errContains != "" && err != nil && !strings.Contains(err.Error(), tc.errContains) {
t.Fatalf("expected error to contain %q, got %q", tc.errContains, err.Error())
}

if !tc.expectErr {
verifyAppliedAssets(t, ctx, kubeClient, tc.assetDir, tc.assetNames)
}
})
}
}

func verifyAppliedAssets(t *testing.T, ctx context.Context, kubeClient *fakecore.Clientset, assetDir string, assetNames []string) {
t.Helper()
for _, name := range assetNames {
switch {
case strings.Contains(name, "_sa.yaml"):
namespace := "openshift-cluster-csi-drivers"
saList, err := kubeClient.CoreV1().ServiceAccounts(namespace).List(ctx, metav1.ListOptions{})
if err != nil {
t.Fatalf("failed to list service accounts: %v", err)
}
if len(saList.Items) == 0 {
t.Fatalf("expected service accounts to be created, but found none")
}
}
}
}
5 changes: 5 additions & 0 deletions pkg/operator/starter.go
Original file line number Diff line number Diff line change
Expand Up @@ -254,6 +254,11 @@ func RunOperator(ctx context.Context, controllerConfig *controllercmd.Controller
c.WaitForCacheSync(ctx)
klog.V(2).Infof("Informers synced")

// Apply some assets earlier. Some controllers won't be happy until they are applied. See, for example, OCPBUGS-99490.
if err = applyPrerequisites(ctx, c.ControlPlaneKubeClient, controllerConfig.EventRecorder, assetDir, opConfig.PrerequisiteAssets); err != nil {
return err
}

if csiOperatorControllerConfig.Precondition != nil {
starterController := StarterController(
c.OperatorClient,
Expand Down