diff --git a/api/v1alpha1/oxidecluster_types.go b/api/v1alpha1/oxidecluster_types.go index e3f3317..492b288 100644 --- a/api/v1alpha1/oxidecluster_types.go +++ b/api/v1alpha1/oxidecluster_types.go @@ -115,6 +115,8 @@ type OxideClusterInitializationStatus struct { // +kubebuilder:printcolumn:name="Age",type="date",JSONPath=".metadata.creationTimestamp" // +kubebuilder:metadata:labels="cluster.x-k8s.io/v1beta1=v1alpha1" // +kubebuilder:metadata:labels="cluster.x-k8s.io/v1beta2=v1alpha1" +// +kubebuilder:metadata:labels="cluster.x-k8s.io/provider=infrastructure-oxide" +// +kubebuilder:metadata:labels="clusterctl.cluster.x-k8s.io=" // OxideCluster is the Schema for the oxideclusters API type OxideCluster struct { metav1.TypeMeta `json:",inline"` diff --git a/api/v1alpha1/oxidemachine_types.go b/api/v1alpha1/oxidemachine_types.go index b5e70df..bbb877c 100644 --- a/api/v1alpha1/oxidemachine_types.go +++ b/api/v1alpha1/oxidemachine_types.go @@ -129,6 +129,8 @@ type OxideMachineStatus struct { // +kubebuilder:printcolumn:name="Age",type="date",JSONPath=".metadata.creationTimestamp" // +kubebuilder:metadata:labels="cluster.x-k8s.io/v1beta1=v1alpha1" // +kubebuilder:metadata:labels="cluster.x-k8s.io/v1beta2=v1alpha1" +// +kubebuilder:metadata:labels="cluster.x-k8s.io/provider=infrastructure-oxide" +// +kubebuilder:metadata:labels="clusterctl.cluster.x-k8s.io=" // OxideMachine is the Schema for the oxidemachines API type OxideMachine struct { metav1.TypeMeta `json:",inline"` diff --git a/api/v1alpha1/oxidemachinetemplate_types.go b/api/v1alpha1/oxidemachinetemplate_types.go index ef1eee1..0dc40d7 100644 --- a/api/v1alpha1/oxidemachinetemplate_types.go +++ b/api/v1alpha1/oxidemachinetemplate_types.go @@ -59,6 +59,8 @@ type OxideMachineTemplateResource struct { // +kubebuilder:subresource:status // +kubebuilder:metadata:labels="cluster.x-k8s.io/v1beta1=v1alpha1" // +kubebuilder:metadata:labels="cluster.x-k8s.io/v1beta2=v1alpha1" +// +kubebuilder:metadata:labels="cluster.x-k8s.io/provider=infrastructure-oxide" +// +kubebuilder:metadata:labels="clusterctl.cluster.x-k8s.io=" // OxideMachineTemplate is the Schema for the oxidemachinetemplates API type OxideMachineTemplate struct { metav1.TypeMeta `json:",inline"` diff --git a/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxideclusters.yaml b/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxideclusters.yaml index 3f63cf0..248fe7e 100644 --- a/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxideclusters.yaml +++ b/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxideclusters.yaml @@ -5,8 +5,10 @@ metadata: annotations: controller-gen.kubebuilder.io/version: v0.21.0 labels: + cluster.x-k8s.io/provider: infrastructure-oxide cluster.x-k8s.io/v1beta1: v1alpha1 cluster.x-k8s.io/v1beta2: v1alpha1 + clusterctl.cluster.x-k8s.io: "" name: oxideclusters.infrastructure.cluster.x-k8s.io spec: group: infrastructure.cluster.x-k8s.io diff --git a/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachines.yaml b/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachines.yaml index 7e90748..14adc45 100644 --- a/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachines.yaml +++ b/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachines.yaml @@ -5,8 +5,10 @@ metadata: annotations: controller-gen.kubebuilder.io/version: v0.21.0 labels: + cluster.x-k8s.io/provider: infrastructure-oxide cluster.x-k8s.io/v1beta1: v1alpha1 cluster.x-k8s.io/v1beta2: v1alpha1 + clusterctl.cluster.x-k8s.io: "" name: oxidemachines.infrastructure.cluster.x-k8s.io spec: group: infrastructure.cluster.x-k8s.io diff --git a/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachinetemplates.yaml b/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachinetemplates.yaml index e8bbcb7..1fd3761 100644 --- a/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachinetemplates.yaml +++ b/charts/cluster-api-provider-oxide/crds/infrastructure.cluster.x-k8s.io_oxidemachinetemplates.yaml @@ -5,8 +5,10 @@ metadata: annotations: controller-gen.kubebuilder.io/version: v0.21.0 labels: + cluster.x-k8s.io/provider: infrastructure-oxide cluster.x-k8s.io/v1beta1: v1alpha1 cluster.x-k8s.io/v1beta2: v1alpha1 + clusterctl.cluster.x-k8s.io: "" name: oxidemachinetemplates.infrastructure.cluster.x-k8s.io spec: group: infrastructure.cluster.x-k8s.io diff --git a/internal/controller/oxidecluster_controller.go b/internal/controller/oxidecluster_controller.go index 43ac99a..ea5c38e 100644 --- a/internal/controller/oxidecluster_controller.go +++ b/internal/controller/oxidecluster_controller.go @@ -29,7 +29,10 @@ import ( "sigs.k8s.io/cluster-api/util" "sigs.k8s.io/cluster-api/util/conditions" "sigs.k8s.io/cluster-api/util/patch" + "sigs.k8s.io/cluster-api/util/paused" + "sigs.k8s.io/cluster-api/util/predicates" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" "sigs.k8s.io/controller-runtime/pkg/handler" @@ -81,6 +84,27 @@ func (r *OxideClusterReconciler) Reconcile( return ctrl.Result{}, err } + cluster, err := util.GetOwnerCluster(ctx, r.Client, oxideCluster.ObjectMeta) + if err != nil { + return ctrl.Result{}, err + } + if cluster == nil { + log.Info("missing ownerRef on OxideCluster", "name", oxideCluster.Name) + return ctrl.Result{}, nil + } + + // Set the Paused condition and return early if paused, e.g. during clusterctl move. + if isPaused, requeue, err := paused.EnsurePausedCondition( + ctx, + r.Client, + cluster, + oxideCluster, + ); err != nil || + isPaused || + requeue { + return ctrl.Result{}, err + } + patchHelper, err := patch.NewHelper(oxideCluster, r.Client) if err != nil { return ctrl.Result{}, fmt.Errorf("building patch helper: %w", err) @@ -91,15 +115,6 @@ func (r *OxideClusterReconciler) Reconcile( } }() - cluster, err := util.GetOwnerCluster(ctx, r.Client, oxideCluster.ObjectMeta) - if err != nil { - return ctrl.Result{}, err - } - if cluster == nil { - log.Info("missing ownerRef on OxideCluster", "name", oxideCluster.Name) - return ctrl.Result{}, nil - } - oxideClient, err := r.OxideClientFactory(ctx, r.Client, oxideCluster) if err != nil { return ctrl.Result{}, err @@ -151,8 +166,23 @@ func (r *OxideClusterReconciler) Reconcile( Reason: infrav1.ReasonFloatingIPProvisioned, }) - // Ensure floating IP is attached to an instance. Use the 0th ready control plane machine if - // unattached. + if err := r.reconcileFloatingIPAttachment(ctx, oxideClient, oxideCluster, ip); err != nil { + return ctrl.Result{}, err + } + + return ctrl.Result{}, nil +} + +// reconcileFloatingIPAttachment ensures the floating IP is attached to a provisioned control plane +// instance, using the 0th ready one if unattached, and sets the FloatingIPAttached condition. +func (r *OxideClusterReconciler) reconcileFloatingIPAttachment( + ctx context.Context, + oxideClient cloud.OxideClient, + oxideCluster *infrav1.OxideCluster, + ip *oxide.FloatingIp, +) error { + log := logf.FromContext(ctx) + shouldAttach := true var machines infrav1.OxideMachineList @@ -167,14 +197,14 @@ func (r *OxideClusterReconciler) Reconcile( clusterv1.MachineControlPlaneLabel: "", }, ); err != nil { - return ctrl.Result{}, fmt.Errorf("listing oxide machines: %w", err) + return fmt.Errorf("listing oxide machines: %w", err) } if ip.InstanceId != "" { for _, machine := range machines.Items { if machine.Spec.ProviderID != "" { instanceID, err := cloud.InstanceIDFromProviderID(machine.Spec.ProviderID) if err != nil { - return ctrl.Result{}, fmt.Errorf("parsing provider id: %w", err) + return fmt.Errorf("parsing provider id: %w", err) } if instanceID == ip.InstanceId { shouldAttach = false @@ -199,11 +229,12 @@ func (r *OxideClusterReconciler) Reconcile( "instance", ip.InstanceId, ) + var err error ip, err = oxideClient.FloatingIpDetach(ctx, oxide.FloatingIpDetachParams{ FloatingIp: oxide.NameOrId(ip.Id), }) if err != nil { - return ctrl.Result{}, fmt.Errorf("detaching floating ip: %w", err) + return fmt.Errorf("detaching floating ip: %w", err) } } @@ -236,7 +267,7 @@ func (r *OxideClusterReconciler) Reconcile( } instanceID, err := cloud.InstanceIDFromProviderID(machine.Spec.ProviderID) if err != nil { - return ctrl.Result{}, fmt.Errorf("parsing provider id: %w", err) + return fmt.Errorf("parsing provider id: %w", err) } log.Info("attaching floating IP", "ip", ip.Ip, "instance", instanceID) ip, err = oxideClient.FloatingIpAttach(ctx, oxide.FloatingIpAttachParams{ @@ -247,7 +278,7 @@ func (r *OxideClusterReconciler) Reconcile( }, }) if err != nil { - return ctrl.Result{}, err + return err } break } @@ -271,7 +302,7 @@ func (r *OxideClusterReconciler) Reconcile( }) } - return ctrl.Result{}, nil + return nil } // floatingIPAllocator builds an oxide.Allocator to provision the floating IP address: @@ -365,11 +396,21 @@ func (r *OxideClusterReconciler) ensureFloatingIPDeleted( // SetupWithManager sets up the controller with the Manager. func (r *OxideClusterReconciler) SetupWithManager(mgr ctrl.Manager) error { + log := mgr.GetLogger().WithValues("controller", "oxidecluster") return ctrl.NewControllerManagedBy(mgr). For(&infrav1.OxideCluster{}). Watches(&infrav1.OxideMachine{}, handler.EnqueueRequestsFromMapFunc( r.oxideMachineToOxideCluster, )). + // Reconcile on Cluster pause transitions, e.g. to resume after a clusterctl move unpauses. + Watches(&clusterv1.Cluster{}, handler.EnqueueRequestsFromMapFunc( + util.ClusterToInfrastructureMapFunc( + context.Background(), + infrav1.GroupVersion.WithKind("OxideCluster"), + mgr.GetClient(), + &infrav1.OxideCluster{}, + ), + ), builder.WithPredicates(predicates.ClusterPausedTransitions(mgr.GetScheme(), log))). Named("oxidecluster"). Complete(r) } diff --git a/internal/controller/oxidecluster_controller_test.go b/internal/controller/oxidecluster_controller_test.go index d7a79ab..1fde482 100644 --- a/internal/controller/oxidecluster_controller_test.go +++ b/internal/controller/oxidecluster_controller_test.go @@ -18,16 +18,25 @@ package controller import ( "context" + "errors" "net/http" "testing" - "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" "go.uber.org/mock/gomock" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + clusterv1 "sigs.k8s.io/cluster-api/api/core/v1beta2" + "sigs.k8s.io/cluster-api/util/conditions" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" infrav1 "github.com/oxidecomputer/cluster-api-provider-oxide/api/v1alpha1" + "github.com/oxidecomputer/cluster-api-provider-oxide/internal/cloud" "github.com/oxidecomputer/cluster-api-provider-oxide/internal/cloud/mock" "github.com/oxidecomputer/oxide.go/oxide" - clusterv1 "sigs.k8s.io/cluster-api/api/core/v1beta2" ) // httpErr constructs an *oxide.HTTPError with a stub HTTPResponse so that its @@ -39,6 +48,93 @@ func httpErr(code string) *oxide.HTTPError { } } +// newPauseTestScheme builds a scheme with the CAPI and Oxide types registered. +func newPauseTestScheme(t *testing.T) *runtime.Scheme { + t.Helper() + scheme := runtime.NewScheme() + require.NoError(t, clusterv1.AddToScheme(scheme)) + require.NoError(t, infrav1.AddToScheme(scheme)) + return scheme +} + +// getPausedCondition re-fetches obj and returns its Paused condition, or nil if unset. +func getPausedCondition(t *testing.T, c client.Client, obj interface { + client.Object + conditions.Getter +}) *metav1.Condition { + t.Helper() + if err := c.Get(context.Background(), client.ObjectKeyFromObject(obj), obj); err != nil { + t.Fatalf("getting %T: %v", obj, err) + } + return conditions.Get(obj, clusterv1.PausedCondition) +} + +func TestOxideClusterReconcilePaused(t *testing.T) { + scheme := newPauseTestScheme(t) + cluster := &clusterv1.Cluster{ + ObjectMeta: metav1.ObjectMeta{Name: "test", Namespace: "default"}, + Spec: clusterv1.ClusterSpec{Paused: new(true)}, + } + oxideCluster := &infrav1.OxideCluster{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test", + Namespace: "default", + OwnerReferences: []metav1.OwnerReference{{ + APIVersion: clusterv1.GroupVersion.String(), + Kind: "Cluster", + Name: "test", + UID: "test-uid", + }}, + }, + } + k8sClient := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(cluster, oxideCluster). + WithStatusSubresource(&infrav1.OxideCluster{}). + Build() + + factoryCalls := 0 + r := &OxideClusterReconciler{ + Client: k8sClient, + Scheme: scheme, + OxideClientFactory: func(context.Context, client.Client, *infrav1.OxideCluster) (cloud.OxideClient, error) { + factoryCalls++ + return nil, errors.New("halting test reconcile") + }, + } + ctx := context.Background() + req := ctrl.Request{NamespacedName: types.NamespacedName{Namespace: "default", Name: "test"}} + + // While paused, the first reconcile sets the Paused condition and requeues, and subsequent + // reconciles skip. The Oxide client must never be constructed. + for range 2 { + result, err := r.Reconcile(ctx, req) + require.NoError(t, err) + require.Equal(t, ctrl.Result{}, result) + } + require.Equal(t, 0, factoryCalls) + cond := getPausedCondition(t, k8sClient, oxideCluster) + require.NotNil(t, cond) + require.Equal(t, metav1.ConditionTrue, cond.Status) + + // Unpause the Cluster. The next reconcile only flips the Paused condition; the one after + // resumes normal reconciliation and constructs the Oxide client. + require.NoError(t, k8sClient.Get(ctx, client.ObjectKeyFromObject(cluster), cluster)) + cluster.Spec.Paused = new(false) + require.NoError(t, k8sClient.Update(ctx, cluster)) + + _, err := r.Reconcile(ctx, req) + require.NoError(t, err) + require.Equal(t, 0, factoryCalls) + cond = getPausedCondition(t, k8sClient, oxideCluster) + require.NotNil(t, cond) + require.Equal(t, metav1.ConditionFalse, cond.Status) + + _, err = r.Reconcile(ctx, req) + require.ErrorContains(t, err, "halting test reconcile") + require.Equal(t, 1, factoryCalls) +} + func TestEnsureFloatingIPExists(t *testing.T) { wantIP := &oxide.FloatingIp{ Id: "ip-id", @@ -96,11 +192,11 @@ func TestEnsureFloatingIPExists(t *testing.T) { "ip-name", ) if tc.wantErr != "" { - assert.ErrorContains(t, gotErr, tc.wantErr) - assert.Nil(t, gotIP) + require.ErrorContains(t, gotErr, tc.wantErr) + require.Nil(t, gotIP) } else { - assert.NoError(t, gotErr) - assert.Equal(t, wantIP, gotIP) + require.NoError(t, gotErr) + require.Equal(t, wantIP, gotIP) } }) } @@ -145,9 +241,9 @@ func TestEnsureFloatingIPDeleted(t *testing.T) { "ip-name", ) if tc.wantErr != "" { - assert.ErrorContains(t, gotErr, tc.wantErr) + require.ErrorContains(t, gotErr, tc.wantErr) } else { - assert.NoError(t, gotErr) + require.NoError(t, gotErr) } }) } @@ -224,7 +320,7 @@ func TestFloatingIPAllocator(t *testing.T) { } { t.Run(tc.name, func(t *testing.T) { got := floatingIPAllocator(tc.cluster) - assert.Equal(t, tc.want, got) + require.Equal(t, tc.want, got) }) } } diff --git a/internal/controller/oxidemachine_controller.go b/internal/controller/oxidemachine_controller.go index 3a93858..9f073f9 100644 --- a/internal/controller/oxidemachine_controller.go +++ b/internal/controller/oxidemachine_controller.go @@ -32,7 +32,10 @@ import ( "sigs.k8s.io/cluster-api/util" "sigs.k8s.io/cluster-api/util/conditions" "sigs.k8s.io/cluster-api/util/patch" + "sigs.k8s.io/cluster-api/util/paused" + "sigs.k8s.io/cluster-api/util/predicates" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" "sigs.k8s.io/controller-runtime/pkg/handler" @@ -77,6 +80,32 @@ func (r *OxideMachineReconciler) Reconcile( return ctrl.Result{}, err } + machine, err := util.GetOwnerMachine(ctx, r.Client, oxideMachine.ObjectMeta) + if err != nil { + return ctrl.Result{}, err + } + if machine == nil { + log.Info("waiting for owner machine reference") + return ctrl.Result{}, nil + } + + cluster, err := util.GetClusterFromMetadata(ctx, r.Client, machine.ObjectMeta) + if err != nil { + return ctrl.Result{}, fmt.Errorf("getting owner cluster: %w", err) + } + + // Set the Paused condition and return early if paused, e.g. during clusterctl move. + if isPaused, requeue, err := paused.EnsurePausedCondition( + ctx, + r.Client, + cluster, + oxideMachine, + ); err != nil || + isPaused || + requeue { + return ctrl.Result{}, err + } + patchHelper, err := patch.NewHelper(oxideMachine, r.Client) if err != nil { return ctrl.Result{}, fmt.Errorf("building patch helper: %w", err) @@ -87,15 +116,6 @@ func (r *OxideMachineReconciler) Reconcile( } }() - machine, err := util.GetOwnerMachine(ctx, r.Client, oxideMachine.ObjectMeta) - if err != nil { - return ctrl.Result{}, err - } - if machine == nil { - log.Info("waiting for owner machine reference") - return ctrl.Result{}, nil - } - clusterName := machine.Labels[clusterv1.ClusterNameLabel] oxideCluster := &infrav1.OxideCluster{} @@ -114,7 +134,6 @@ func (r *OxideMachineReconciler) Reconcile( projectName := oxideCluster.Spec.Project instanceName := getInstanceName(oxideMachine) bootDiskName := getBootDiskName(oxideMachine) - nicName := getNicName(oxideMachine) if !oxideMachine.DeletionTimestamp.IsZero() { return r.handleDelete( @@ -131,103 +150,13 @@ func (r *OxideMachineReconciler) Reconcile( // Ensure instance exists. Instance creation idempotently creates the disk and NIC as well, so // create all resources in a single request. - var instance *oxide.Instance - if oxideMachine.Spec.ProviderID == "" { - // Fetch the UserData from the bootstrap secret. If the secret isn't set on the spec yet, - // mark the OxideMachine as unready, and wait for an update to DataSecretName to trigger a - // new reconcile. - bootstrapSecretName := machine.Spec.Bootstrap.DataSecretName - if bootstrapSecretName == nil { - conditions.Set(oxideMachine, metav1.Condition{ - Type: clusterv1.ReadyCondition, - Status: metav1.ConditionFalse, - Reason: clusterv1.WaitingForBootstrapDataReason, - }) - return ctrl.Result{}, nil - } - var bootstrapSecret corev1.Secret - if err := r.Get(ctx, client.ObjectKey{ - Namespace: machine.Namespace, - Name: *bootstrapSecretName, - }, &bootstrapSecret); err != nil { - return ctrl.Result{}, fmt.Errorf("fetching bootstrap secret: %w", err) - } - if _, ok := bootstrapSecret.Data["value"]; !ok { - return ctrl.Result{}, fmt.Errorf( - "missing `value` key in bootstrap secret %s", - *bootstrapSecretName, - ) - } - - instance, err = oxideClient.InstanceCreate(ctx, oxide.InstanceCreateParams{ - Project: oxide.NameOrId(projectName), - Body: &oxide.InstanceCreate{ - Name: oxide.Name(instanceName), - Hostname: oxide.Hostname(instanceName), - Ncpus: oxide.InstanceCpuCount(oxideMachine.Spec.NCpus), - Memory: oxide.ByteCount(oxideMachine.Spec.Memory.Value()), - Start: new(true), - AntiAffinityGroups: toNamesOrIds(oxideMachine.Spec.AntiAffinityGroups), - SshPublicKeys: toNamesOrIds(oxideMachine.Spec.SSHPublicKeys), - UserData: base64.StdEncoding.EncodeToString( - bootstrapSecret.Data["value"], - ), - BootDisk: oxide.InstanceDiskAttachment{ - Value: oxide.InstanceDiskAttachmentCreate{ - Name: oxide.Name(bootDiskName), - Size: oxide.ByteCount(oxideMachine.Spec.DiskSize.Value()), - DiskBackend: oxide.DiskBackend{ - Value: oxide.DiskBackendDistributed{ - DiskSource: oxide.DiskSource{ - Value: oxide.DiskSourceImage{ - ImageId: oxideMachine.Spec.ImageID, - }, - }, - }, - }, - }, - }, - Disks: disksFromOxideMachine(oxideMachine), - NetworkInterfaces: oxide.InstanceNetworkInterfaceAttachment{ - Value: oxide.InstanceNetworkInterfaceAttachmentCreate{ - Params: []oxide.InstanceNetworkInterfaceCreate{ - { - Name: oxide.Name(nicName), - VpcName: oxide.Name(oxideCluster.Spec.VPC), - SubnetName: oxide.Name(oxideCluster.Spec.Subnet), - }, - }, - }, - }, - ExternalIps: externalIPsFromMachine(oxideMachine), - }, - }) - if err != nil { - // Look up the instance if creation failed with a conflict, and adopt the existing - // instance if found. Note: if an instance was created out of band with unexpected - // parameters, it will be adopted as well; operators shouldn't create or modify these - // instances outside the reconciler. - if !errors.Is(err, oxide.ErrObjectAlreadyExists) { - return ctrl.Result{}, fmt.Errorf("creating oxide instance: %w", err) - } - instance, err = oxideClient.InstanceView(ctx, oxide.InstanceViewParams{ - Project: oxide.NameOrId(projectName), - Instance: oxide.NameOrId(instanceName), - }) - if err != nil { - return ctrl.Result{}, fmt.Errorf("viewing existing oxide instance: %w", err) - } - } - - oxideMachine.Spec.ProviderID = cloud.NewProviderID(instance.Id) - } else { - instance, err = oxideClient.InstanceView(ctx, oxide.InstanceViewParams{ - Project: oxide.NameOrId(projectName), - Instance: oxide.NameOrId(instanceName), - }) - if err != nil { - return ctrl.Result{}, fmt.Errorf("fetching oxide instance: %w", err) - } + instance, err := r.ensureInstance(ctx, oxideClient, oxideMachine, machine, oxideCluster) + if err != nil { + return ctrl.Result{}, err + } + if instance == nil { + // Waiting for bootstrap data; a Machine update triggers the next reconcile. + return ctrl.Result{}, nil } instanceRunning, instance, err := r.ensureInstanceRunning(ctx, oxideClient, instance) @@ -290,6 +219,119 @@ func (r *OxideMachineReconciler) Reconcile( return ctrl.Result{}, nil } +// ensureInstance idempotently creates or adopts the Oxide instance for the machine, setting +// Spec.ProviderID on creation. It returns a nil instance while waiting for bootstrap data. +func (r *OxideMachineReconciler) ensureInstance( + ctx context.Context, + oxideClient cloud.OxideClient, + oxideMachine *infrav1.OxideMachine, + machine *clusterv1.Machine, + oxideCluster *infrav1.OxideCluster, +) (*oxide.Instance, error) { + projectName := oxideCluster.Spec.Project + instanceName := getInstanceName(oxideMachine) + + if oxideMachine.Spec.ProviderID != "" { + instance, err := oxideClient.InstanceView(ctx, oxide.InstanceViewParams{ + Project: oxide.NameOrId(projectName), + Instance: oxide.NameOrId(instanceName), + }) + if err != nil { + return nil, fmt.Errorf("fetching oxide instance: %w", err) + } + return instance, nil + } + + // Fetch the UserData from the bootstrap secret. If the secret isn't set on the spec yet, + // mark the OxideMachine as unready, and wait for an update to DataSecretName to trigger a + // new reconcile. + bootstrapSecretName := machine.Spec.Bootstrap.DataSecretName + if bootstrapSecretName == nil { + conditions.Set(oxideMachine, metav1.Condition{ + Type: clusterv1.ReadyCondition, + Status: metav1.ConditionFalse, + Reason: clusterv1.WaitingForBootstrapDataReason, + }) + return nil, nil + } + var bootstrapSecret corev1.Secret + if err := r.Get(ctx, client.ObjectKey{ + Namespace: machine.Namespace, + Name: *bootstrapSecretName, + }, &bootstrapSecret); err != nil { + return nil, fmt.Errorf("fetching bootstrap secret: %w", err) + } + if _, ok := bootstrapSecret.Data["value"]; !ok { + return nil, fmt.Errorf( + "missing `value` key in bootstrap secret %s", + *bootstrapSecretName, + ) + } + + instance, err := oxideClient.InstanceCreate(ctx, oxide.InstanceCreateParams{ + Project: oxide.NameOrId(projectName), + Body: &oxide.InstanceCreate{ + Name: oxide.Name(instanceName), + Hostname: oxide.Hostname(instanceName), + Ncpus: oxide.InstanceCpuCount(oxideMachine.Spec.NCpus), + Memory: oxide.ByteCount(oxideMachine.Spec.Memory.Value()), + Start: new(true), + AntiAffinityGroups: toNamesOrIds(oxideMachine.Spec.AntiAffinityGroups), + SshPublicKeys: toNamesOrIds(oxideMachine.Spec.SSHPublicKeys), + UserData: base64.StdEncoding.EncodeToString( + bootstrapSecret.Data["value"], + ), + BootDisk: oxide.InstanceDiskAttachment{ + Value: oxide.InstanceDiskAttachmentCreate{ + Name: oxide.Name(getBootDiskName(oxideMachine)), + Size: oxide.ByteCount(oxideMachine.Spec.DiskSize.Value()), + DiskBackend: oxide.DiskBackend{ + Value: oxide.DiskBackendDistributed{ + DiskSource: oxide.DiskSource{ + Value: oxide.DiskSourceImage{ + ImageId: oxideMachine.Spec.ImageID, + }, + }, + }, + }, + }, + }, + Disks: disksFromOxideMachine(oxideMachine), + NetworkInterfaces: oxide.InstanceNetworkInterfaceAttachment{ + Value: oxide.InstanceNetworkInterfaceAttachmentCreate{ + Params: []oxide.InstanceNetworkInterfaceCreate{ + { + Name: oxide.Name(getNicName(oxideMachine)), + VpcName: oxide.Name(oxideCluster.Spec.VPC), + SubnetName: oxide.Name(oxideCluster.Spec.Subnet), + }, + }, + }, + }, + ExternalIps: externalIPsFromMachine(oxideMachine), + }, + }) + if err != nil { + // Look up the instance if creation failed with a conflict, and adopt the existing + // instance if found. Note: if an instance was created out of band with unexpected + // parameters, it will be adopted as well; operators shouldn't create or modify these + // instances outside the reconciler. + if !errors.Is(err, oxide.ErrObjectAlreadyExists) { + return nil, fmt.Errorf("creating oxide instance: %w", err) + } + instance, err = oxideClient.InstanceView(ctx, oxide.InstanceViewParams{ + Project: oxide.NameOrId(projectName), + Instance: oxide.NameOrId(instanceName), + }) + if err != nil { + return nil, fmt.Errorf("viewing existing oxide instance: %w", err) + } + } + + oxideMachine.Spec.ProviderID = cloud.NewProviderID(instance.Id) + return instance, nil +} + // ensureInstanceRunning ensures that the given instance is running, starting it if necessary. // // * If stopped, start and requeue. @@ -487,6 +529,15 @@ func (r *OxideMachineReconciler) ensureDiskDeleted( // SetupWithManager sets up the controller with the Manager. func (r *OxideMachineReconciler) SetupWithManager(mgr ctrl.Manager) error { + log := mgr.GetLogger().WithValues("controller", "oxidemachine") + clusterToOxideMachines, err := util.ClusterToTypedObjectsMapper( + mgr.GetClient(), + &infrav1.OxideMachineList{}, + mgr.GetScheme(), + ) + if err != nil { + return fmt.Errorf("creating cluster to oxidemachines mapper: %w", err) + } return ctrl.NewControllerManagedBy(mgr). For(&infrav1.OxideMachine{}). // The OxideMachine reconciler depends on the state of the parent Machine; watch the parent @@ -495,6 +546,13 @@ func (r *OxideMachineReconciler) SetupWithManager(mgr ctrl.Manager) error { &clusterv1.Machine{}, handler.EnqueueRequestsFromMapFunc(util.MachineToInfrastructureMapFunc(infrav1.GroupVersion.WithKind("OxideMachine"))), ). + + // Reconcile on Cluster pause transitions, e.g. to resume after a clusterctl move unpauses. + Watches( + &clusterv1.Cluster{}, + handler.EnqueueRequestsFromMapFunc(clusterToOxideMachines), + builder.WithPredicates(predicates.ClusterPausedTransitions(mgr.GetScheme(), log)), + ). Named("oxidemachine"). Complete(r) } diff --git a/internal/controller/oxidemachine_controller_test.go b/internal/controller/oxidemachine_controller_test.go index a3bd058..14821d5 100644 --- a/internal/controller/oxidemachine_controller_test.go +++ b/internal/controller/oxidemachine_controller_test.go @@ -18,19 +18,103 @@ package controller import ( "context" + "errors" "testing" "github.com/stretchr/testify/assert" "go.uber.org/mock/gomock" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" clusterv1 "sigs.k8s.io/cluster-api/api/core/v1beta2" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" infrav1 "github.com/oxidecomputer/cluster-api-provider-oxide/api/v1alpha1" + "github.com/oxidecomputer/cluster-api-provider-oxide/internal/cloud" "github.com/oxidecomputer/cluster-api-provider-oxide/internal/cloud/mock" "github.com/oxidecomputer/oxide.go/oxide" ) +func TestOxideMachineReconcilePaused(t *testing.T) { + scheme := newPauseTestScheme(t) + cluster := &clusterv1.Cluster{ + ObjectMeta: metav1.ObjectMeta{Name: "test", Namespace: "default"}, + Spec: clusterv1.ClusterSpec{Paused: new(true)}, + } + oxideCluster := &infrav1.OxideCluster{ + ObjectMeta: metav1.ObjectMeta{Name: "test", Namespace: "default"}, + } + machine := &clusterv1.Machine{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-machine", + Namespace: "default", + Labels: map[string]string{clusterv1.ClusterNameLabel: "test"}, + }, + } + oxideMachine := &infrav1.OxideMachine{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-machine", + Namespace: "default", + OwnerReferences: []metav1.OwnerReference{{ + APIVersion: clusterv1.GroupVersion.String(), + Kind: "Machine", + Name: "test-machine", + UID: "test-machine-uid", + }}, + }, + } + k8sClient := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(cluster, oxideCluster, machine, oxideMachine). + WithStatusSubresource(&infrav1.OxideMachine{}). + Build() + + factoryCalls := 0 + r := &OxideMachineReconciler{ + Client: k8sClient, + Scheme: scheme, + OxideClientFactory: func(context.Context, client.Client, *infrav1.OxideCluster) (cloud.OxideClient, error) { + factoryCalls++ + return nil, errors.New("halting test reconcile") + }, + } + ctx := context.Background() + req := ctrl.Request{ + NamespacedName: types.NamespacedName{Namespace: "default", Name: "test-machine"}, + } + + // While paused, the first reconcile sets the Paused condition and requeues, and subsequent + // reconciles skip. The Oxide client must never be constructed. + for range 2 { + result, err := r.Reconcile(ctx, req) + assert.NoError(t, err) + assert.Equal(t, ctrl.Result{}, result) + } + assert.Equal(t, 0, factoryCalls) + if cond := getPausedCondition(t, k8sClient, oxideMachine); assert.NotNil(t, cond) { + assert.Equal(t, metav1.ConditionTrue, cond.Status) + } + + // Unpause the Cluster. The next reconcile only flips the Paused condition; the one after + // resumes normal reconciliation and constructs the Oxide client. + assert.NoError(t, k8sClient.Get(ctx, client.ObjectKeyFromObject(cluster), cluster)) + cluster.Spec.Paused = new(false) + assert.NoError(t, k8sClient.Update(ctx, cluster)) + + _, err := r.Reconcile(ctx, req) + assert.NoError(t, err) + assert.Equal(t, 0, factoryCalls) + if cond := getPausedCondition(t, k8sClient, oxideMachine); assert.NotNil(t, cond) { + assert.Equal(t, metav1.ConditionFalse, cond.Status) + } + + _, err = r.Reconcile(ctx, req) + assert.ErrorContains(t, err, "halting test reconcile") + assert.Equal(t, 1, factoryCalls) +} + func TestEnsureInstanceRunning(t *testing.T) { for _, tc := range []struct { name string