diff --git a/charts/k3k/templates/crds/k3k.io_clusters.yaml b/charts/k3k/templates/crds/k3k.io_clusters.yaml
index 29543207..c0853c6d 100644
--- a/charts/k3k/templates/crds/k3k.io_clusters.yaml
+++ b/charts/k3k/templates/crds/k3k.io_clusters.yaml
@@ -810,6 +810,25 @@ spec:
required:
- enabled
type: object
+ storageClasses:
+ default:
+ enabled: false
+ description: StorageClasses resources sync configuration.
+ properties:
+ enabled:
+ default: false
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
type: object
tlsSANs:
description: TLSSANs specifies subject alternative names for the K3s
@@ -955,6 +974,141 @@ spec:
description: priorityClass is the priority class enforced by the
active VirtualClusterPolicy.
type: string
+ sync:
+ description: sync is the SyncConfig enforced by the active VirtualClusterPolicy.
+ properties:
+ configMaps:
+ default:
+ enabled: true
+ description: ConfigMaps resources sync configuration.
+ properties:
+ enabled:
+ default: true
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
+ ingresses:
+ default:
+ enabled: false
+ description: Ingresses resources sync configuration.
+ properties:
+ enabled:
+ default: false
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
+ persistentVolumeClaims:
+ default:
+ enabled: true
+ description: PersistentVolumeClaims resources sync configuration.
+ properties:
+ enabled:
+ default: true
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
+ priorityClasses:
+ default:
+ enabled: false
+ description: PriorityClasses resources sync configuration.
+ properties:
+ enabled:
+ default: false
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
+ secrets:
+ default:
+ enabled: true
+ description: Secrets resources sync configuration.
+ properties:
+ enabled:
+ default: true
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ type: object
+ services:
+ default:
+ enabled: true
+ description: Services resources sync configuration.
+ properties:
+ enabled:
+ default: true
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
+ storageClasses:
+ default:
+ enabled: false
+ description: StorageClasses resources sync configuration.
+ properties:
+ enabled:
+ default: false
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
+ type: object
required:
- name
type: object
diff --git a/charts/k3k/templates/crds/k3k.io_virtualclusterpolicies.yaml b/charts/k3k/templates/crds/k3k.io_virtualclusterpolicies.yaml
index 28a8790f..2371ea92 100644
--- a/charts/k3k/templates/crds/k3k.io_virtualclusterpolicies.yaml
+++ b/charts/k3k/templates/crds/k3k.io_virtualclusterpolicies.yaml
@@ -343,6 +343,25 @@ spec:
required:
- enabled
type: object
+ storageClasses:
+ default:
+ enabled: false
+ description: StorageClasses resources sync configuration.
+ properties:
+ enabled:
+ default: false
+ description: Enabled is an on/off switch for syncing resources.
+ type: boolean
+ selector:
+ additionalProperties:
+ type: string
+ description: |-
+ Selector specifies set of labels of the resources that will be synced, if empty
+ then all resources of the given type will be synced.
+ type: object
+ required:
+ - enabled
+ type: object
type: object
type: object
status:
diff --git a/docs/crds/crds.adoc b/docs/crds/crds.adoc
index c60d56d3..dd90713f 100644
--- a/docs/crds/crds.adoc
+++ b/docs/crds/crds.adoc
@@ -61,6 +61,7 @@ _Appears In:_
| *`priorityClass`* __string__ | priorityClass is the priority class enforced by the active VirtualClusterPolicy. + | |
| *`nodeSelector`* __object (keys:string, values:string)__ | nodeSelector is a node selector enforced by the active VirtualClusterPolicy. + | |
+| *`sync`* __xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-syncconfig[$$SyncConfig$$]__ | sync is the SyncConfig enforced by the active VirtualClusterPolicy. + | |
|===
@@ -217,7 +218,7 @@ Each entry defines a secret and its mount path within the pods. + | |
-ConfigMapSyncConfig specifies the sync options for services.
+ConfigMapSyncConfig specifies the sync options for ConfigMaps.
@@ -352,7 +353,7 @@ _Appears In:_
-IngressSyncConfig specifies the sync options for services.
+IngressSyncConfig specifies the sync options for Ingresses.
@@ -463,7 +464,7 @@ _Appears In:_
-PersistentVolumeClaimSyncConfig specifies the sync options for services.
+PersistentVolumeClaimSyncConfig specifies the sync options for PersistentVolumeClaims.
@@ -501,7 +502,7 @@ _Appears In:_
-PriorityClassSyncConfig specifies the sync options for services.
+PriorityClassSyncConfig specifies the sync options for PriorityClasses.
@@ -568,7 +569,7 @@ This can be 'server', 'agent', or 'all' (for both). + | | Enum: [server agent a
-SecretSyncConfig specifies the sync options for services.
+SecretSyncConfig specifies the sync options for Secrets.
@@ -590,7 +591,7 @@ then all resources of the given type will be synced. + | |
-ServiceSyncConfig specifies the sync options for services.
+ServiceSyncConfig specifies the sync options for Services.
@@ -607,6 +608,28 @@ then all resources of the given type will be synced. + | |
|===
+[id="{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-storageclasssyncconfig"]
+=== StorageClassSyncConfig
+
+
+
+StorageClassSyncConfig specifies the sync options for StorageClasses.
+
+
+
+_Appears In:_
+
+* xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-syncconfig[$$SyncConfig$$]
+
+[cols="25a,55a,10a,10a", options="header"]
+|===
+| Field | Description | Default | Validation
+| *`enabled`* __boolean__ | Enabled is an on/off switch for syncing resources. + | false |
+| *`selector`* __object (keys:string, values:string)__ | Selector specifies set of labels of the resources that will be synced, if empty +
+then all resources of the given type will be synced. + | |
+|===
+
+
[id="{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-syncconfig"]
=== SyncConfig
@@ -618,6 +641,7 @@ SyncConfig will contain the resources that should be synced from virtual cluster
_Appears In:_
+* xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-appliedpolicy[$$AppliedPolicy$$]
* xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-clusterspec[$$ClusterSpec$$]
* xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-virtualclusterpolicyspec[$$VirtualClusterPolicySpec$$]
@@ -630,6 +654,7 @@ _Appears In:_
| *`ingresses`* __xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-ingresssyncconfig[$$IngressSyncConfig$$]__ | Ingresses resources sync configuration. + | { enabled:false } |
| *`persistentVolumeClaims`* __xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-persistentvolumeclaimsyncconfig[$$PersistentVolumeClaimSyncConfig$$]__ | PersistentVolumeClaims resources sync configuration. + | { enabled:true } |
| *`priorityClasses`* __xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-priorityclasssyncconfig[$$PriorityClassSyncConfig$$]__ | PriorityClasses resources sync configuration. + | { enabled:false } |
+| *`storageClasses`* __xref:{anchor_prefix}-github-com-rancher-k3k-pkg-apis-k3k-io-v1beta1-storageclasssyncconfig[$$StorageClassSyncConfig$$]__ | StorageClasses resources sync configuration. + | { enabled:false } |
|===
diff --git a/docs/crds/crds.md b/docs/crds/crds.md
index a44a1ade..5a0a0c7e 100644
--- a/docs/crds/crds.md
+++ b/docs/crds/crds.md
@@ -48,6 +48,7 @@ _Appears in:_
| `name` _string_ | name is the name of the VirtualClusterPolicy currently applied to this cluster. | | MinLength: 1
|
| `priorityClass` _string_ | priorityClass is the priority class enforced by the active VirtualClusterPolicy. | | |
| `nodeSelector` _object (keys:string, values:string)_ | nodeSelector is a node selector enforced by the active VirtualClusterPolicy. | | |
+| `sync` _[SyncConfig](#syncconfig)_ | sync is the SyncConfig enforced by the active VirtualClusterPolicy. | | |
#### Cluster
@@ -162,7 +163,7 @@ _Appears in:_
-ConfigMapSyncConfig specifies the sync options for services.
+ConfigMapSyncConfig specifies the sync options for ConfigMaps.
@@ -270,7 +271,7 @@ _Appears in:_
-IngressSyncConfig specifies the sync options for services.
+IngressSyncConfig specifies the sync options for Ingresses.
@@ -352,7 +353,7 @@ _Appears in:_
-PersistentVolumeClaimSyncConfig specifies the sync options for services.
+PersistentVolumeClaimSyncConfig specifies the sync options for PersistentVolumeClaims.
@@ -383,7 +384,7 @@ _Appears in:_
-PriorityClassSyncConfig specifies the sync options for services.
+PriorityClassSyncConfig specifies the sync options for PriorityClasses.
@@ -423,7 +424,7 @@ _Appears in:_
-SecretSyncConfig specifies the sync options for services.
+SecretSyncConfig specifies the sync options for Secrets.
@@ -440,7 +441,7 @@ _Appears in:_
-ServiceSyncConfig specifies the sync options for services.
+ServiceSyncConfig specifies the sync options for Services.
@@ -453,6 +454,23 @@ _Appears in:_
| `selector` _object (keys:string, values:string)_ | Selector specifies set of labels of the resources that will be synced, if empty
then all resources of the given type will be synced. | | |
+#### StorageClassSyncConfig
+
+
+
+StorageClassSyncConfig specifies the sync options for StorageClasses.
+
+
+
+_Appears in:_
+- [SyncConfig](#syncconfig)
+
+| Field | Description | Default | Validation |
+| --- | --- | --- | --- |
+| `enabled` _boolean_ | Enabled is an on/off switch for syncing resources. | false | |
+| `selector` _object (keys:string, values:string)_ | Selector specifies set of labels of the resources that will be synced, if empty
then all resources of the given type will be synced. | | |
+
+
#### SyncConfig
@@ -462,6 +480,7 @@ SyncConfig will contain the resources that should be synced from virtual cluster
_Appears in:_
+- [AppliedPolicy](#appliedpolicy)
- [ClusterSpec](#clusterspec)
- [VirtualClusterPolicySpec](#virtualclusterpolicyspec)
@@ -473,6 +492,7 @@ _Appears in:_
| `ingresses` _[IngressSyncConfig](#ingresssyncconfig)_ | Ingresses resources sync configuration. | \{ enabled:false \} | |
| `persistentVolumeClaims` _[PersistentVolumeClaimSyncConfig](#persistentvolumeclaimsyncconfig)_ | PersistentVolumeClaims resources sync configuration. | \{ enabled:true \} | |
| `priorityClasses` _[PriorityClassSyncConfig](#priorityclasssyncconfig)_ | PriorityClasses resources sync configuration. | \{ enabled:false \} | |
+| `storageClasses` _[StorageClassSyncConfig](#storageclasssyncconfig)_ | StorageClasses resources sync configuration. | \{ enabled:false \} | |
#### VirtualClusterPolicy
diff --git a/pkg/apis/k3k.io/v1beta1/types.go b/pkg/apis/k3k.io/v1beta1/types.go
index e2e019ec..d72eff68 100644
--- a/pkg/apis/k3k.io/v1beta1/types.go
+++ b/pkg/apis/k3k.io/v1beta1/types.go
@@ -249,9 +249,14 @@ type SyncConfig struct {
// +kubebuilder:default={"enabled": false}
// +optional
PriorityClasses PriorityClassSyncConfig `json:"priorityClasses"`
+ // StorageClasses resources sync configuration.
+ //
+ // +kubebuilder:default={"enabled": false}
+ // +optional
+ StorageClasses StorageClassSyncConfig `json:"storageClasses"`
}
-// SecretSyncConfig specifies the sync options for services.
+// SecretSyncConfig specifies the sync options for Secrets.
type SecretSyncConfig struct {
// Enabled is an on/off switch for syncing resources.
//
@@ -266,7 +271,7 @@ type SecretSyncConfig struct {
Selector map[string]string `json:"selector,omitempty"`
}
-// ServiceSyncConfig specifies the sync options for services.
+// ServiceSyncConfig specifies the sync options for Services.
type ServiceSyncConfig struct {
// Enabled is an on/off switch for syncing resources.
//
@@ -281,7 +286,7 @@ type ServiceSyncConfig struct {
Selector map[string]string `json:"selector,omitempty"`
}
-// ConfigMapSyncConfig specifies the sync options for services.
+// ConfigMapSyncConfig specifies the sync options for ConfigMaps.
type ConfigMapSyncConfig struct {
// Enabled is an on/off switch for syncing resources.
//
@@ -296,7 +301,7 @@ type ConfigMapSyncConfig struct {
Selector map[string]string `json:"selector,omitempty"`
}
-// IngressSyncConfig specifies the sync options for services.
+// IngressSyncConfig specifies the sync options for Ingresses.
type IngressSyncConfig struct {
// Enabled is an on/off switch for syncing resources.
//
@@ -311,7 +316,7 @@ type IngressSyncConfig struct {
Selector map[string]string `json:"selector,omitempty"`
}
-// PersistentVolumeClaimSyncConfig specifies the sync options for services.
+// PersistentVolumeClaimSyncConfig specifies the sync options for PersistentVolumeClaims.
type PersistentVolumeClaimSyncConfig struct {
// Enabled is an on/off switch for syncing resources.
//
@@ -326,7 +331,7 @@ type PersistentVolumeClaimSyncConfig struct {
Selector map[string]string `json:"selector,omitempty"`
}
-// PriorityClassSyncConfig specifies the sync options for services.
+// PriorityClassSyncConfig specifies the sync options for PriorityClasses.
type PriorityClassSyncConfig struct {
// Enabled is an on/off switch for syncing resources.
//
@@ -341,6 +346,21 @@ type PriorityClassSyncConfig struct {
Selector map[string]string `json:"selector,omitempty"`
}
+// StorageClassSyncConfig specifies the sync options for StorageClasses.
+type StorageClassSyncConfig struct {
+ // Enabled is an on/off switch for syncing resources.
+ //
+ // +kubebuilder:default=false
+ // +required
+ Enabled bool `json:"enabled"`
+
+ // Selector specifies set of labels of the resources that will be synced, if empty
+ // then all resources of the given type will be synced.
+ //
+ // +optional
+ Selector map[string]string `json:"selector,omitempty"`
+}
+
// ClusterMode is the possible provisioning mode of a Cluster.
//
// +kubebuilder:validation:Enum=shared;virtual
@@ -584,6 +604,11 @@ type AppliedPolicy struct {
//
// +optional
NodeSelector map[string]string `json:"nodeSelector,omitempty"`
+
+ // sync is the SyncConfig enforced by the active VirtualClusterPolicy.
+ //
+ // +optional
+ Sync *SyncConfig `json:"sync,omitempty"`
}
// ClusterPhase is a high-level summary of the cluster's current lifecycle state.
diff --git a/pkg/apis/k3k.io/v1beta1/zz_generated.deepcopy.go b/pkg/apis/k3k.io/v1beta1/zz_generated.deepcopy.go
index b4b2d14a..eab97897 100644
--- a/pkg/apis/k3k.io/v1beta1/zz_generated.deepcopy.go
+++ b/pkg/apis/k3k.io/v1beta1/zz_generated.deepcopy.go
@@ -40,6 +40,11 @@ func (in *AppliedPolicy) DeepCopyInto(out *AppliedPolicy) {
(*out)[key] = val
}
}
+ if in.Sync != nil {
+ in, out := &in.Sync, &out.Sync
+ *out = new(SyncConfig)
+ (*in).DeepCopyInto(*out)
+ }
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new AppliedPolicy.
@@ -578,6 +583,28 @@ func (in *ServiceSyncConfig) DeepCopy() *ServiceSyncConfig {
return out
}
+// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
+func (in *StorageClassSyncConfig) DeepCopyInto(out *StorageClassSyncConfig) {
+ *out = *in
+ if in.Selector != nil {
+ in, out := &in.Selector, &out.Selector
+ *out = make(map[string]string, len(*in))
+ for key, val := range *in {
+ (*out)[key] = val
+ }
+ }
+}
+
+// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new StorageClassSyncConfig.
+func (in *StorageClassSyncConfig) DeepCopy() *StorageClassSyncConfig {
+ if in == nil {
+ return nil
+ }
+ out := new(StorageClassSyncConfig)
+ in.DeepCopyInto(out)
+ return out
+}
+
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
func (in *SyncConfig) DeepCopyInto(out *SyncConfig) {
*out = *in
@@ -587,6 +614,7 @@ func (in *SyncConfig) DeepCopyInto(out *SyncConfig) {
in.Ingresses.DeepCopyInto(&out.Ingresses)
in.PersistentVolumeClaims.DeepCopyInto(&out.PersistentVolumeClaims)
in.PriorityClasses.DeepCopyInto(&out.PriorityClasses)
+ in.StorageClasses.DeepCopyInto(&out.StorageClasses)
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new SyncConfig.
diff --git a/pkg/controller/cluster/cluster.go b/pkg/controller/cluster/cluster.go
index 74957212..ddc802e7 100644
--- a/pkg/controller/cluster/cluster.go
+++ b/pkg/controller/cluster/cluster.go
@@ -11,6 +11,7 @@ import (
"k8s.io/apimachinery/pkg/api/equality"
"k8s.io/apimachinery/pkg/api/meta"
+ "k8s.io/apimachinery/pkg/labels"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/discovery"
@@ -28,6 +29,7 @@ import (
v1 "k8s.io/api/core/v1"
networkingv1 "k8s.io/api/networking/v1"
rbacv1 "k8s.io/api/rbac/v1"
+ storagev1 "k8s.io/api/storage/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
ctrl "sigs.k8s.io/controller-runtime"
@@ -47,11 +49,18 @@ const (
clusterFinalizerName = "cluster.k3k.io/finalizer"
ClusterInvalidName = "system"
+ SyncEnabledLabelKey = "k3k.io/sync-enabled"
+ SyncSourceLabelKey = "k3k.io/sync-source"
+ SyncSourceHostLabel = "host"
+
defaultVirtualClusterCIDR = "10.52.0.0/16"
defaultVirtualServiceCIDR = "10.53.0.0/16"
defaultSharedClusterCIDR = "10.42.0.0/16"
defaultSharedServiceCIDR = "10.43.0.0/16"
memberRemovalTimeout = time.Minute * 1
+
+ storageClassEnabledIndexField = "spec.sync.storageClasses.enabled"
+ storageClassStatusEnabledIndexField = "status.policy.sync.storageClasses.enabled"
)
var (
@@ -115,15 +124,82 @@ func Add(ctx context.Context, mgr manager.Manager, config *Config, maxConcurrent
},
}
+ // index the 'spec.sync.storageClasses.enabled' field
+ err = mgr.GetCache().IndexField(ctx, &v1beta1.Cluster{}, storageClassEnabledIndexField, func(rawObj client.Object) []string {
+ vc := rawObj.(*v1beta1.Cluster)
+
+ if vc.Spec.Sync != nil && vc.Spec.Sync.StorageClasses.Enabled {
+ return []string{"true"}
+ }
+
+ return []string{"false"}
+ })
+ if err != nil {
+ return err
+ }
+
+ // index the 'status.policy.sync.storageClasses.enabled' field
+ err = mgr.GetCache().IndexField(ctx, &v1beta1.Cluster{}, storageClassStatusEnabledIndexField, func(rawObj client.Object) []string {
+ vc := rawObj.(*v1beta1.Cluster)
+
+ if vc.Status.Policy != nil && vc.Status.Policy.Sync != nil && vc.Status.Policy.Sync.StorageClasses.Enabled {
+ return []string{"true"}
+ }
+
+ return []string{"false"}
+ })
+ if err != nil {
+ return err
+ }
+
return ctrl.NewControllerManagedBy(mgr).
For(&v1beta1.Cluster{}).
Watches(&v1.Namespace{}, namespaceEventHandler(&reconciler)).
+ Watches(&storagev1.StorageClass{},
+ handler.EnqueueRequestsFromMapFunc(reconciler.mapStorageClassToCluster),
+ ).
Owns(&apps.StatefulSet{}).
Owns(&v1.Service{}).
WithOptions(ctrlcontroller.Options{MaxConcurrentReconciles: maxConcurrentReconciles}).
Complete(&reconciler)
}
+func (r *ClusterReconciler) mapStorageClassToCluster(ctx context.Context, obj client.Object) []reconcile.Request {
+ log := ctrl.LoggerFrom(ctx)
+
+ if _, ok := obj.(*storagev1.StorageClass); !ok {
+ return nil
+ }
+
+ // Merge and deduplicate clusters
+ allClusters := make(map[types.NamespacedName]struct{})
+
+ var specClusterList v1beta1.ClusterList
+ if err := r.Client.List(ctx, &specClusterList, client.MatchingFields{storageClassEnabledIndexField: "true"}); err != nil {
+ log.Error(err, "error listing clusters with spec sync enabled for storageclass sync")
+ } else {
+ for _, cluster := range specClusterList.Items {
+ allClusters[client.ObjectKeyFromObject(&cluster)] = struct{}{}
+ }
+ }
+
+ var statusClusterList v1beta1.ClusterList
+ if err := r.Client.List(ctx, &statusClusterList, client.MatchingFields{storageClassStatusEnabledIndexField: "true"}); err != nil {
+ log.Error(err, "error listing clusters with status sync enabled for storageclass sync")
+ } else {
+ for _, cluster := range statusClusterList.Items {
+ allClusters[client.ObjectKeyFromObject(&cluster)] = struct{}{}
+ }
+ }
+
+ requests := make([]reconcile.Request, 0, len(allClusters))
+ for key := range allClusters {
+ requests = append(requests, reconcile.Request{NamespacedName: key})
+ }
+
+ return requests
+}
+
func namespaceEventHandler(r *ClusterReconciler) handler.Funcs {
return handler.Funcs{
// We don't need to update for create or delete events
@@ -350,11 +426,22 @@ func (c *ClusterReconciler) reconcile(ctx context.Context, cluster *v1beta1.Clus
return err
}
+ if err := c.bindClusterRoles(ctx, cluster); err != nil {
+ return err
+ }
+
if err := c.ensureKubeconfigSecret(ctx, cluster, serviceIP, 443); err != nil {
return err
}
- return c.bindClusterRoles(ctx, cluster)
+ // Important: if you need to call the Server API of the Virtual Cluster
+ // this needs to be done AFTER he kubeconfig has been generated
+
+ if err := c.ensureStorageClasses(ctx, cluster); err != nil {
+ return err
+ }
+
+ return nil
}
// ensureBootstrapSecret will create or update the Secret containing the bootstrap data from the k3s server
@@ -620,6 +707,120 @@ func (c *ClusterReconciler) ensureIngress(ctx context.Context, cluster *v1beta1.
return nil
}
+func (c *ClusterReconciler) ensureStorageClasses(ctx context.Context, cluster *v1beta1.Cluster) error {
+ log := ctrl.LoggerFrom(ctx)
+ log.V(1).Info("Ensuring cluster StorageClasses")
+
+ virtualClient, err := newVirtualClient(ctx, c.Client, cluster.Name, cluster.Namespace)
+ if err != nil {
+ return fmt.Errorf("failed creating virtual client: %w", err)
+ }
+
+ appliedSync := cluster.Spec.Sync.DeepCopy()
+
+ // If a policy is applied to the virtual cluster we need to use its SyncConfig, if available
+ if cluster.Status.Policy != nil && cluster.Status.Policy.Sync != nil {
+ appliedSync = cluster.Status.Policy.Sync
+ }
+
+ // If storageclass sync is disabled, clean up any managed storage classes.
+ if appliedSync == nil || !appliedSync.StorageClasses.Enabled {
+ err := virtualClient.DeleteAllOf(ctx, &storagev1.StorageClass{}, client.MatchingLabels{SyncSourceLabelKey: SyncSourceHostLabel})
+ return client.IgnoreNotFound(err)
+ }
+
+ var hostStorageClasses storagev1.StorageClassList
+ if err := c.Client.List(ctx, &hostStorageClasses); err != nil {
+ return fmt.Errorf("failed listing host storageclasses: %w", err)
+ }
+
+ // filter the StorageClasses disabled for the sync, and the one not matching the selector
+ filteredHostStorageClasses := make(map[string]storagev1.StorageClass)
+
+ for _, sc := range hostStorageClasses.Items {
+ syncEnabled, found := sc.Labels[SyncEnabledLabelKey]
+
+ // if sync is disabled -> continue
+ if found && syncEnabled != "true" {
+ log.V(1).Info("sync is disabled", "sc-name", sc.Name)
+ continue
+ }
+
+ // if selector doesn't match -> continue
+ // an empty selector matche everything
+ selector := labels.SelectorFromSet(appliedSync.StorageClasses.Selector)
+ if !selector.Matches(labels.Set(sc.Labels)) {
+ log.V(1).Info("selector not matching", "sc-name", sc.Name)
+ continue
+ }
+
+ log.V(1).Info("keeping storageclass", "sc-name", sc.Name)
+
+ filteredHostStorageClasses[sc.Name] = sc
+ }
+
+ var virtStorageClasses storagev1.StorageClassList
+ if err = virtualClient.List(ctx, &virtStorageClasses, client.MatchingLabels{SyncSourceLabelKey: SyncSourceHostLabel}); err != nil {
+ return fmt.Errorf("failed listing virtual storageclasses: %w", err)
+ }
+
+ // delete StorageClasses with the sync disabled
+
+ for _, sc := range virtStorageClasses.Items {
+ if _, found := filteredHostStorageClasses[sc.Name]; !found {
+ log.V(1).Info("deleting storageclass", "sc-name", sc.Name)
+
+ if errDelete := virtualClient.Delete(ctx, &sc); errDelete != nil {
+ log.Error(errDelete, "failed to delete virtual storageclass", "name", sc.Name)
+ err = errors.Join(err, errDelete)
+ }
+ }
+ }
+
+ for _, hostSc := range filteredHostStorageClasses {
+ log.V(1).Info("updating storageclass", "sc-name", hostSc.Name)
+
+ virtualSc := hostSc.DeepCopy()
+
+ virtualSc.ObjectMeta = metav1.ObjectMeta{
+ Name: hostSc.Name,
+ Labels: hostSc.Labels,
+ Annotations: hostSc.Annotations,
+ }
+
+ _, errCreateOrUpdate := controllerutil.CreateOrUpdate(ctx, virtualClient, virtualSc, func() error {
+ virtualSc.Annotations = hostSc.Annotations
+
+ virtualSc.Labels = hostSc.Labels
+ if len(virtualSc.Labels) == 0 {
+ virtualSc.Labels = make(map[string]string)
+ }
+
+ virtualSc.Labels[SyncSourceLabelKey] = SyncSourceHostLabel
+
+ virtualSc.Provisioner = hostSc.Provisioner
+ virtualSc.Parameters = hostSc.Parameters
+ virtualSc.ReclaimPolicy = hostSc.ReclaimPolicy
+ virtualSc.MountOptions = hostSc.MountOptions
+ virtualSc.AllowVolumeExpansion = hostSc.AllowVolumeExpansion
+ virtualSc.VolumeBindingMode = hostSc.VolumeBindingMode
+ virtualSc.AllowedTopologies = hostSc.AllowedTopologies
+
+ return nil
+ })
+ if errCreateOrUpdate != nil {
+ log.Error(errCreateOrUpdate, "failed to create or update virtual storageclass", "name", virtualSc.Name)
+ err = errors.Join(err, errCreateOrUpdate)
+ }
+ }
+
+ if err != nil {
+ return fmt.Errorf("failed to sync storageclasses: %w", err)
+ }
+
+ return nil
+}
+
func (c *ClusterReconciler) server(ctx context.Context, cluster *v1beta1.Cluster, server *server.Server) error {
log := ctrl.LoggerFrom(ctx)
@@ -742,11 +943,6 @@ func (c *ClusterReconciler) validate(cluster *v1beta1.Cluster, policy v1beta1.Vi
}
}
- // validate sync policy
- if !equality.Semantic.DeepEqual(cluster.Spec.Sync, policy.Spec.Sync) {
- return fmt.Errorf("sync configuration %v is not allowed by the policy %q", cluster.Spec.Sync, policy.Name)
- }
-
return nil
}
diff --git a/pkg/controller/policy/policy.go b/pkg/controller/policy/policy.go
index f3d8222c..bde4fffd 100644
--- a/pkg/controller/policy/policy.go
+++ b/pkg/controller/policy/policy.go
@@ -476,6 +476,7 @@ func (c *VirtualClusterPolicyReconciler) reconcileClusters(ctx context.Context,
Name: policy.Name,
PriorityClass: &policy.Spec.DefaultPriorityClass,
NodeSelector: policy.Spec.DefaultNodeSelector,
+ Sync: policy.Spec.Sync,
}
if !reflect.DeepEqual(origStatus, &cluster.Status) {
diff --git a/tests/cluster_policy_sync_storageclass_test.go b/tests/cluster_policy_sync_storageclass_test.go
new file mode 100644
index 00000000..fac20b67
--- /dev/null
+++ b/tests/cluster_policy_sync_storageclass_test.go
@@ -0,0 +1,199 @@
+package k3k_test
+
+import (
+ "context"
+ "time"
+
+ "k8s.io/utils/ptr"
+ "sigs.k8s.io/controller-runtime/pkg/client"
+
+ corev1 "k8s.io/api/core/v1"
+ storagev1 "k8s.io/api/storage/v1"
+ apierrors "k8s.io/apimachinery/pkg/api/errors"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+
+ "github.com/rancher/k3k/pkg/apis/k3k.io/v1beta1"
+ "github.com/rancher/k3k/pkg/controller/cluster"
+ "github.com/rancher/k3k/pkg/controller/policy"
+
+ . "github.com/onsi/ginkgo/v2"
+ . "github.com/onsi/gomega"
+)
+
+var _ = When("a shared mode cluster is created in a namespace with a policy", Ordered, Label(e2eTestLabel), func() {
+ var (
+ ctx context.Context
+ virtualCluster *VirtualCluster
+ vcp *v1beta1.VirtualClusterPolicy
+ )
+
+ BeforeAll(func() {
+ ctx = context.Background()
+
+ // 1. Create StorageClasses in host
+ storageClassEnabled := &storagev1.StorageClass{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "sc-policy-enabled-",
+ Labels: map[string]string{
+ cluster.SyncEnabledLabelKey: "true",
+ },
+ },
+ Provisioner: "my-provisioner",
+ }
+
+ var err error
+ storageClassEnabled, err = k8s.StorageV1().StorageClasses().Create(ctx, storageClassEnabled, metav1.CreateOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ storageClassDisabled := &storagev1.StorageClass{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "sc-policy-disabled-",
+ Labels: map[string]string{
+ cluster.SyncEnabledLabelKey: "false",
+ },
+ },
+ Provisioner: "my-provisioner",
+ }
+
+ storageClassDisabled, err = k8s.StorageV1().StorageClasses().Create(ctx, storageClassDisabled, metav1.CreateOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ // 2. Create VirtualClusterPolicy with StorageClass sync enabled
+ vcp = &v1beta1.VirtualClusterPolicy{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "vcp-sync-sc-",
+ },
+ Spec: v1beta1.VirtualClusterPolicySpec{
+ Sync: &v1beta1.SyncConfig{
+ StorageClasses: v1beta1.StorageClassSyncConfig{
+ Enabled: true,
+ },
+ },
+ },
+ }
+ err = k8sClient.Create(ctx, vcp)
+ Expect(err).To(Not(HaveOccurred()))
+
+ // 3. Create Namespace with policy label
+ ns := &corev1.Namespace{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "ns-vcp-",
+ Labels: map[string]string{
+ policy.PolicyNameLabelKey: vcp.Name,
+ },
+ },
+ }
+ // We use the k8s clientset for namespace creation to stay consistent with other tests
+ ns, err = k8s.CoreV1().Namespaces().Create(ctx, ns, metav1.CreateOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ // 4. Create VirtualCluster in that namespace
+ // The cluster doesn't have storage class sync enabled in its spec
+ clusterObj := NewCluster(ns.Name)
+ clusterObj.Spec.Sync = &v1beta1.SyncConfig{
+ StorageClasses: v1beta1.StorageClassSyncConfig{
+ Enabled: false,
+ },
+ }
+ clusterObj.Spec.Expose.NodePort.ServerPort = ptr.To[int32](30000)
+
+ CreateCluster(clusterObj)
+
+ client, restConfig, kubeconfig := NewVirtualK8sClientAndKubeconfig(clusterObj)
+ virtualCluster = &VirtualCluster{
+ Cluster: clusterObj,
+ RestConfig: restConfig,
+ Client: client,
+ Kubeconfig: kubeconfig,
+ }
+
+ DeferCleanup(func() {
+ DeleteNamespaces(ns.Name)
+
+ err = k8s.StorageV1().StorageClasses().Delete(ctx, storageClassEnabled.Name, metav1.DeleteOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ err = k8s.StorageV1().StorageClasses().Delete(ctx, storageClassDisabled.Name, metav1.DeleteOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ err = k8sClient.Delete(ctx, vcp)
+ Expect(err).To(Not(HaveOccurred()))
+ })
+ })
+
+ It("has the storage classes sync enabled from the policy", func() {
+ Eventually(func(g Gomega) {
+ key := client.ObjectKeyFromObject(virtualCluster.Cluster)
+ g.Expect(k8sClient.Get(ctx, key, virtualCluster.Cluster)).To(Succeed())
+ g.Expect(virtualCluster.Cluster.Status.Policy).To(Not(BeNil()))
+ g.Expect(virtualCluster.Cluster.Status.Policy.Sync).To(Not(BeNil()))
+ g.Expect(virtualCluster.Cluster.Status.Policy.Sync.StorageClasses.Enabled).To(BeTrue())
+ }).
+ WithTimeout(time.Second * 30).
+ WithPolling(time.Second).
+ Should(Succeed())
+ })
+
+ It("will sync host storage classes with the sync enabled in the host", func() {
+ Eventually(func(g Gomega) {
+ hostStorageClasses, err := k8s.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ g.Expect(err).To(Not(HaveOccurred()))
+
+ for _, hostSC := range hostStorageClasses.Items {
+ // We only care about the storage classes we created for this test to avoid noise
+ if hostSC.Labels[cluster.SyncEnabledLabelKey] == "true" {
+ _, err := virtualCluster.Client.StorageV1().StorageClasses().Get(ctx, hostSC.Name, metav1.GetOptions{})
+ g.Expect(err).To(Not(HaveOccurred()))
+ }
+ }
+ }).
+ WithPolling(time.Second).
+ WithTimeout(time.Second * 60).
+ Should(Succeed())
+ })
+
+ It("will not sync host storage classes with the sync disabled in the host", func() {
+ Eventually(func(g Gomega) {
+ hostStorageClasses, err := k8s.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ g.Expect(err).To(Not(HaveOccurred()))
+
+ for _, hostSC := range hostStorageClasses.Items {
+ if hostSC.Labels[cluster.SyncEnabledLabelKey] == "false" {
+ _, err := virtualCluster.Client.StorageV1().StorageClasses().Get(ctx, hostSC.Name, metav1.GetOptions{})
+ g.Expect(err).To(HaveOccurred())
+ g.Expect(apierrors.IsNotFound(err)).To(BeTrue())
+ }
+ }
+ }).
+ WithPolling(time.Second).
+ WithTimeout(time.Second * 60).
+ Should(Succeed())
+ })
+
+ When("disabling the storage class sync in the policy", Ordered, func() {
+ BeforeAll(func() {
+ original := vcp.DeepCopy()
+ vcp.Spec.Sync.StorageClasses.Enabled = false
+ err := k8sClient.Patch(ctx, vcp, client.MergeFrom(original))
+ Expect(err).To(Not(HaveOccurred()))
+ })
+
+ It("will remove the synced storage classes from the virtual cluster", func() {
+ Eventually(func(g Gomega) {
+ hostStorageClasses, err := k8s.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ g.Expect(err).To(Not(HaveOccurred()))
+
+ for _, hostSC := range hostStorageClasses.Items {
+ if hostSC.Labels[cluster.SyncEnabledLabelKey] == "true" {
+ _, err := virtualCluster.Client.StorageV1().StorageClasses().Get(ctx, hostSC.Name, metav1.GetOptions{})
+ g.Expect(err).To(HaveOccurred())
+ g.Expect(apierrors.IsNotFound(err)).To(BeTrue())
+ }
+ }
+ }).
+ WithPolling(time.Second).
+ WithTimeout(time.Second * 60).
+ Should(Succeed())
+ })
+ })
+})
diff --git a/tests/cluster_sync_storageclass_test.go b/tests/cluster_sync_storageclass_test.go
new file mode 100644
index 00000000..99df140c
--- /dev/null
+++ b/tests/cluster_sync_storageclass_test.go
@@ -0,0 +1,194 @@
+package k3k_test
+
+import (
+ "context"
+ "time"
+
+ "sigs.k8s.io/controller-runtime/pkg/client"
+
+ storagev1 "k8s.io/api/storage/v1"
+ apierrors "k8s.io/apimachinery/pkg/api/errors"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+
+ "github.com/rancher/k3k/pkg/controller/cluster"
+
+ . "github.com/onsi/ginkgo/v2"
+ . "github.com/onsi/gomega"
+)
+
+var _ = When("a shared mode cluster is created", Ordered, Label(e2eTestLabel), func() {
+ var (
+ ctx context.Context
+ virtualCluster *VirtualCluster
+ )
+
+ BeforeAll(func() {
+ ctx = context.Background()
+ virtualCluster = NewVirtualCluster()
+
+ storageClassEnabled := &storagev1.StorageClass{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "sc-",
+ Labels: map[string]string{
+ cluster.SyncEnabledLabelKey: "true",
+ },
+ },
+ Provisioner: "my-provisioner",
+ }
+
+ storageClassEnabled, err := k8s.StorageV1().StorageClasses().Create(ctx, storageClassEnabled, metav1.CreateOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ storageClassDisabled := &storagev1.StorageClass{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "sc-",
+ Labels: map[string]string{
+ cluster.SyncEnabledLabelKey: "false",
+ },
+ },
+ Provisioner: "my-provisioner",
+ }
+
+ storageClassDisabled, err = k8s.StorageV1().StorageClasses().Create(ctx, storageClassDisabled, metav1.CreateOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ DeferCleanup(func() {
+ DeleteNamespaces(virtualCluster.Cluster.Namespace)
+
+ err = k8s.StorageV1().StorageClasses().Delete(ctx, storageClassEnabled.Name, metav1.DeleteOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ err = k8s.StorageV1().StorageClasses().Delete(ctx, storageClassDisabled.Name, metav1.DeleteOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+ })
+ })
+
+ It("has disabled the storage classes sync", func() {
+ Expect(virtualCluster.Cluster.Spec.Sync).To(Not(BeNil()))
+ Expect(virtualCluster.Cluster.Spec.Sync.StorageClasses.Enabled).To(BeFalse())
+ })
+
+ It("doesn't have storage classes", func() {
+ virtualStorageClasses, err := virtualCluster.Client.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+ Expect(virtualStorageClasses.Items).To(HaveLen(0))
+ })
+
+ It("has some storage classes in the host", func() {
+ hostStorageClasses, err := k8s.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+ Expect(hostStorageClasses.Items).To(Not(HaveLen(0)))
+ })
+
+ It("can create storage classes in the virtual cluster", func() {
+ storageClass := &storagev1.StorageClass{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "sc-",
+ },
+ Provisioner: "my-provisioner",
+ }
+
+ storageClass, err := virtualCluster.Client.StorageV1().StorageClasses().Create(ctx, storageClass, metav1.CreateOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ virtualStorageClasses, err := virtualCluster.Client.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+ Expect(virtualStorageClasses.Items).To(HaveLen(1))
+ Expect(virtualStorageClasses.Items[0].Name).To(Equal(storageClass.Name))
+ })
+
+ When("enabling the storage class sync", Ordered, func() {
+ BeforeAll(func() {
+ GinkgoWriter.Println("Enabling the storage class sync")
+
+ original := virtualCluster.Cluster.DeepCopy()
+
+ virtualCluster.Cluster.Spec.Sync.StorageClasses.Enabled = true
+
+ err := k8sClient.Patch(ctx, virtualCluster.Cluster, client.MergeFrom(original))
+ Expect(err).To(Not(HaveOccurred()))
+
+ Eventually(func(g Gomega) {
+ key := client.ObjectKeyFromObject(virtualCluster.Cluster)
+ g.Expect(k8sClient.Get(ctx, key, virtualCluster.Cluster)).To(Succeed())
+ g.Expect(virtualCluster.Cluster.Spec.Sync.StorageClasses.Enabled).To(BeTrue())
+ }).
+ WithTimeout(time.Second * 10).
+ WithPolling(time.Second).
+ Should(Succeed())
+ })
+
+ It("will sync host storage classes with the sync enabled", func() {
+ Eventually(func(g Gomega) {
+ hostStorageClasses, err := k8s.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ for _, hostSC := range hostStorageClasses.Items {
+ _, err := virtualCluster.Client.StorageV1().StorageClasses().Get(ctx, hostSC.Name, metav1.GetOptions{})
+
+ if syncEnabled, found := hostSC.Labels[cluster.SyncEnabledLabelKey]; !found || syncEnabled == "true" {
+ g.Expect(err).To(Not(HaveOccurred()))
+ }
+ }
+ }).
+ MustPassRepeatedly(5).
+ WithPolling(time.Second).
+ WithTimeout(time.Second * 30).
+ Should(Succeed())
+ })
+
+ It("will not sync host storage classes with the sync disabled", func() {
+ Eventually(func(g Gomega) {
+ hostStorageClasses, err := k8s.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ for _, hostSC := range hostStorageClasses.Items {
+ _, err := virtualCluster.Client.StorageV1().StorageClasses().Get(ctx, hostSC.Name, metav1.GetOptions{})
+
+ if hostSC.Labels[cluster.SyncEnabledLabelKey] == "false" {
+ g.Expect(err).To(HaveOccurred())
+ g.Expect(apierrors.IsNotFound(err)).To(BeTrue())
+ }
+ }
+ }).
+ MustPassRepeatedly(5).
+ WithPolling(time.Second).
+ WithTimeout(time.Second * 30).
+ Should(Succeed())
+ })
+ })
+
+ When("editing a synced storage class in the host cluster", Ordered, func() {
+ var syncedStorageClass *storagev1.StorageClass
+
+ BeforeAll(func() {
+ hostStorageClasses, err := k8s.StorageV1().StorageClasses().List(ctx, metav1.ListOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+
+ for _, hostSC := range hostStorageClasses.Items {
+ if syncEnabled, found := hostSC.Labels[cluster.SyncEnabledLabelKey]; !found || syncEnabled == "true" {
+ syncedStorageClass = &hostSC
+ break
+ }
+ }
+
+ Expect(syncedStorageClass).To(Not(BeNil()))
+
+ syncedStorageClass.Labels["foo"] = "bar"
+ _, err = k8s.StorageV1().StorageClasses().Update(ctx, syncedStorageClass, metav1.UpdateOptions{})
+ Expect(err).To(Not(HaveOccurred()))
+ })
+
+ It("will update the synced storage class in the virtual cluster", func() {
+ Eventually(func(g Gomega) {
+ _, err := virtualCluster.Client.StorageV1().StorageClasses().Get(ctx, syncedStorageClass.Name, metav1.GetOptions{})
+ g.Expect(err).To(Not(HaveOccurred()))
+ g.Expect(syncedStorageClass.Labels).Should(HaveKeyWithValue("foo", "bar"))
+ }).
+ MustPassRepeatedly(5).
+ WithPolling(time.Second).
+ WithTimeout(time.Second * 30).
+ Should(Succeed())
+ })
+ })
+})