diff --git a/pkg/controller/osimagestream/osimagestream_controller.go b/pkg/controller/osimagestream/osimagestream_controller.go index 2d9bb4e736..1c31415f8e 100644 --- a/pkg/controller/osimagestream/osimagestream_controller.go +++ b/pkg/controller/osimagestream/osimagestream_controller.go @@ -47,7 +47,7 @@ type Controller struct { eventRecorder record.EventRecorder fgHandler ctrlcommon.FeatureGatesHandler - syncHandler func(key string) error + syncHandler func(ctx context.Context, key string) error osImageStreamLister mcfglistersv1.OSImageStreamLister osImageStreamListerSynced cache.InformerSynced @@ -169,7 +169,7 @@ func (ctrl *Controller) Run(ctx context.Context, workers int) { klog.Info("Starting MachineConfigController-OSImageStreamController") defer klog.Info("Shutting down MachineConfigController-OSImageStreamController") - go wait.Until(ctrl.worker, time.Second, ctx.Done()) + go wait.Until(func() { ctrl.worker(ctx) }, time.Second, ctx.Done()) // Run once at start at least to make sure the OSImageStream is created when creating the cluster ctrl.enqueue() @@ -192,19 +192,19 @@ func (ctrl *Controller) EnsureOSImageStream(ctx context.Context) error { } } -func (ctrl *Controller) worker() { - for ctrl.processNextWorkItem() { +func (ctrl *Controller) worker(ctx context.Context) { + for ctrl.processNextWorkItem(ctx) { } } -func (ctrl *Controller) processNextWorkItem() bool { +func (ctrl *Controller) processNextWorkItem(ctx context.Context) bool { key, quit := ctrl.queue.Get() if quit { return false } defer ctrl.queue.Done(key) - err := ctrl.syncHandler(key) + err := ctrl.syncHandler(ctx, key) ctrl.handleErr(err, key) return true @@ -269,7 +269,7 @@ func (ctrl *Controller) updateOSImageStream(old, cur interface{}) { ctrl.enqueue() } -func (ctrl *Controller) syncOSImageStream(key string) error { +func (ctrl *Controller) syncOSImageStream(ctx context.Context, key string) error { startTime := time.Now() klog.V(4).Infof("Started syncing OSImageStream %q (%v)", key, startTime) defer func() { @@ -286,21 +286,26 @@ func (ctrl *Controller) syncOSImageStream(key string) error { return err } + var current *mcfgv1.OSImageStream if rebuildRequired { - if err := ctrl.buildOSImageStream(context.TODO(), existing); err != nil { + built, err := ctrl.buildOSImageStream(ctx, existing) + if err != nil { return err } + current = built } else if existing != nil { - if err := ctrl.handleOSImageStreamUpdate(context.TODO(), existing); err != nil { + updated, err := ctrl.handleOSImageStreamUpdate(ctx, existing) + if err != nil { return err } + current = updated } ctrl.initialSyncOnce.Do(func() { - if osis, err := ctrl.getExistingOSImageStream(); err == nil && osis != nil { + if current != nil { klog.Infof("OSImageStream synced successfully. Available streams: %s. Default stream: %s", - osimagestream.GetStreamSetsNames(osis.Status.AvailableStreams), - osis.Status.DefaultStream) + osimagestream.GetStreamSetsNames(current.Status.AvailableStreams), + current.Status.DefaultStream) } ctrl.initialSyncDone <- nil }) @@ -353,16 +358,16 @@ func (ctrl *Controller) getExistingOSImageStream() (*mcfgv1.OSImageStream, error return osis, nil } -func (ctrl *Controller) buildOSImageStream(ctx context.Context, existing *mcfgv1.OSImageStream) error { +func (ctrl *Controller) buildOSImageStream(ctx context.Context, existing *mcfgv1.OSImageStream) (*mcfgv1.OSImageStream, error) { klog.Info("Starting building of the OSImageStream instance") clusterVersion, err := osimagestream.GetClusterVersion(ctrl.clusterVersionLister) if err != nil { - return fmt.Errorf("getting cluster version for OSImageStream inspection: %w", err) + return nil, fmt.Errorf("getting cluster version for OSImageStream inspection: %w", err) } image, err := osimagestream.GetReleasePayloadImage(clusterVersion) if err != nil { - return fmt.Errorf("getting the Release Image digest from the ClusterVersion for OSImageStream sync: %w", err) + return nil, fmt.Errorf("getting the Release Image digest from the ClusterVersion for OSImageStream sync: %w", err) } installVersion, err := osimagestream.GetInstallVersion(clusterVersion) @@ -384,7 +389,7 @@ func (ctrl *Controller) buildOSImageStream(ctx context.Context, existing *mcfgv1 InstallVersion: installVersion, }) if err != nil { - return fmt.Errorf("building the OSImageStream: %w", err) + return nil, fmt.Errorf("building the OSImageStream: %w", err) } var updateOSImageStream *mcfgv1.OSImageStream @@ -392,7 +397,7 @@ func (ctrl *Controller) buildOSImageStream(ctx context.Context, existing *mcfgv1 klog.V(4).Info("Creating OSImageStream singleton instance") updateOSImageStream, err = ctrl.mcfgClient.MachineconfigurationV1().OSImageStreams().Create(ctx, osImageStream, metav1.CreateOptions{}) if err != nil { - return fmt.Errorf("error creating the OSImageStream: %w", err) + return nil, fmt.Errorf("error creating the OSImageStream: %w", err) } klog.Infof("Created OSImageStream with %d available streams, default stream: %s", len(osImageStream.Status.AvailableStreams), osImageStream.Status.DefaultStream) @@ -406,7 +411,7 @@ func (ctrl *Controller) buildOSImageStream(ctx context.Context, existing *mcfgv1 maps.Copy(desired.Annotations, osImageStream.Annotations) updateOSImageStream, err = ctrl.mcfgClient.MachineconfigurationV1().OSImageStreams().Update(ctx, desired, metav1.UpdateOptions{}) if err != nil { - return fmt.Errorf("error updating the OSImageStream: %w", err) + return nil, fmt.Errorf("error updating the OSImageStream: %w", err) } } @@ -415,36 +420,38 @@ func (ctrl *Controller) buildOSImageStream(ctx context.Context, existing *mcfgv1 MachineconfigurationV1(). OSImageStreams(). UpdateStatus(ctx, updateOSImageStream, metav1.UpdateOptions{}); err != nil { - return fmt.Errorf("error updating the OSImageStream status: %w", err) + return nil, fmt.Errorf("error updating the OSImageStream status: %w", err) } - return nil + return updateOSImageStream, nil } -func (ctrl *Controller) handleOSImageStreamUpdate(ctx context.Context, existing *mcfgv1.OSImageStream) error { +func (ctrl *Controller) handleOSImageStreamUpdate(ctx context.Context, existing *mcfgv1.OSImageStream) (*mcfgv1.OSImageStream, error) { requestedDefault := osimagestream.GetOSImageStreamSpecDefault(existing) if requestedDefault == "" { - return nil + return existing, nil } currentDefault := existing.Status.DefaultStream if currentDefault != requestedDefault { if _, err := osimagestream.GetOSImageStreamSetByName(existing, requestedDefault); err != nil { - return fmt.Errorf("syncing default OSImageStream with OSImageStream %s: %w", requestedDefault, err) + return nil, fmt.Errorf("syncing default OSImageStream with OSImageStream %s: %w", requestedDefault, err) } osis := existing.DeepCopy() osis.Status.DefaultStream = requestedDefault - if _, err := ctrl.mcfgClient. + updated, err := ctrl.mcfgClient. MachineconfigurationV1(). OSImageStreams(). - UpdateStatus(ctx, osis, metav1.UpdateOptions{}); err != nil { - return fmt.Errorf("updating the default OSImageStream status: %w", err) + UpdateStatus(ctx, osis, metav1.UpdateOptions{}) + if err != nil { + return nil, fmt.Errorf("updating the default OSImageStream status: %w", err) } klog.Infof("OSImageStream default has changed from %s to %s", currentDefault, requestedDefault) + return updated, nil } - return nil + return existing, nil } // osImageStreamRequiresRebuild checks if the OSImageStream needs to be created or updated. @@ -466,7 +473,6 @@ func osImageStreamRequiresRebuild(osImageStream *mcfgv1.OSImageStream, releaseIm return !ok || storedVersion != version.Hash } - // Retain implements CacheEvicter — it keeps cached digests referenced by the // OSImage and OSExtensionsImage of each available stream in the OSImageStream. func (ctrl *Controller) Retain(digests []string) []string {