Skip to content
Open
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
62 changes: 34 additions & 28 deletions pkg/controller/osimagestream/osimagestream_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand All @@ -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
Expand Down Expand Up @@ -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() {
Expand All @@ -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
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

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
})
Expand Down Expand Up @@ -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)
Expand All @@ -384,15 +389,15 @@ 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
if existing == nil {
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)
Expand All @@ -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)
}
}

Expand All @@ -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.
Expand All @@ -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 {
Expand Down