• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

heathcliff26 / kube-upgrade / 19393334557

15 Nov 2025 05:41PM UTC coverage: 72.307% (+0.2%) from 72.155%
19393334557

Pull #192

github

web-flow
Merge bcbbff0ce into 9625339c7
Pull Request #192: upgrade-controller: Use validating webhook to ensure only a single plan exists

18 of 28 new or added lines in 3 files covered. (64.29%)

1094 of 1513 relevant lines covered (72.31%)

11.69 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

52.02
/pkg/upgrade-controller/controller/controller.go
1
package controller
2

3
import (
4
        "context"
5
        "fmt"
6
        "log/slog"
7
        "time"
8

9
        api "github.com/heathcliff26/kube-upgrade/pkg/apis/kubeupgrade/v1alpha3"
10
        "github.com/heathcliff26/kube-upgrade/pkg/constants"
11
        "golang.org/x/mod/semver"
12
        appv1 "k8s.io/api/apps/v1"
13
        corev1 "k8s.io/api/core/v1"
14
        "k8s.io/client-go/rest"
15
        ctrl "sigs.k8s.io/controller-runtime"
16
        "sigs.k8s.io/controller-runtime/pkg/cache"
17
        "sigs.k8s.io/controller-runtime/pkg/client"
18
        "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
19
        "sigs.k8s.io/controller-runtime/pkg/healthz"
20
        "sigs.k8s.io/controller-runtime/pkg/manager"
21
        "sigs.k8s.io/controller-runtime/pkg/manager/signals"
22
)
23

24
const (
25
        defaultUpgradedImage = "ghcr.io/heathcliff26/kube-upgraded"
26
        upgradedImageEnv     = "UPGRADED_IMAGE"
27
        upgradedTagEnv       = "UPGRADED_TAG"
28
)
29

30
type controller struct {
31
        client.Client
32
        manager       manager.Manager
33
        namespace     string
34
        upgradedImage string
35
}
36

37
// Run make generate when changing these comments
38
// +kubebuilder:rbac:groups=kubeupgrade.heathcliff.eu,resources=kubeupgradeplans,verbs=get;list;watch;create;update;patch;delete
39
// +kubebuilder:rbac:groups=kubeupgrade.heathcliff.eu,resources=kubeupgradeplans/status,verbs=get;update;patch
40
// +kubebuilder:rbac:groups="",resources=nodes,verbs=list;watch;update
41
// +kubebuilder:rbac:groups="",namespace=kube-upgrade,resources=events,verbs=create;patch
42
// +kubebuilder:rbac:groups="coordination.k8s.io",namespace=kube-upgrade,resources=leases,verbs=create;get;update
43
// +kubebuilder:rbac:groups="apps",namespace=kube-upgrade,resources=daemonsets,verbs=list;watch;create;update;delete
44
// +kubebuilder:rbac:groups="",namespace=kube-upgrade,resources=configmaps,verbs=list;watch;create;update;delete
45

46
func NewController(name string) (*controller, error) {
1✔
47
        config, err := rest.InClusterConfig()
1✔
48
        if err != nil {
2✔
49
                return nil, err
1✔
50
        }
1✔
51

52
        ns, err := GetNamespace()
×
53
        if err != nil {
×
54
                return nil, err
×
55
        }
×
56

NEW
57
        scheme, err := newScheme()
×
58
        if err != nil {
×
59
                return nil, err
×
60
        }
×
61

62
        mgr, err := ctrl.NewManager(config, manager.Options{
×
63
                Scheme:                        scheme,
×
64
                LeaderElection:                true,
×
65
                LeaderElectionNamespace:       ns,
×
66
                LeaderElectionID:              name,
×
67
                LeaderElectionReleaseOnCancel: true,
×
68
                LeaseDuration:                 Pointer(time.Minute),
×
69
                RenewDeadline:                 Pointer(10 * time.Second),
×
70
                RetryPeriod:                   Pointer(5 * time.Second),
×
71
                HealthProbeBindAddress:        ":9090",
×
72
                Cache: cache.Options{
×
73
                        DefaultNamespaces: map[string]cache.Config{ns: {}},
×
74
                },
×
75
        })
×
76
        if err != nil {
×
77
                return nil, err
×
78
        }
×
79
        err = mgr.AddHealthzCheck("healthz", healthz.Ping)
×
80
        if err != nil {
×
81
                return nil, err
×
82
        }
×
83
        err = mgr.AddReadyzCheck("readyz", healthz.Ping)
×
84
        if err != nil {
×
85
                return nil, err
×
86
        }
×
87

88
        return &controller{
×
89
                Client:        mgr.GetClient(),
×
90
                manager:       mgr,
×
91
                namespace:     ns,
×
92
                upgradedImage: GetUpgradedImage(),
×
93
        }, nil
×
94
}
95

96
func (c *controller) Run() error {
×
97
        err := ctrl.NewControllerManagedBy(c.manager).
×
98
                For(&api.KubeUpgradePlan{}).
×
99
                Owns(&appv1.DaemonSet{}).
×
100
                Owns(&corev1.ConfigMap{}).
×
101
                Complete(c)
×
102
        if err != nil {
×
103
                return err
×
104
        }
×
105

106
        err = ctrl.NewWebhookManagedBy(c.manager).
×
107
                For(&api.KubeUpgradePlan{}).
×
108
                WithDefaulter(&planMutatingHook{}).
×
NEW
109
                WithValidator(&planValidatingHook{
×
NEW
110
                        Client: c.Client,
×
NEW
111
                }).
×
112
                Complete()
×
113
        if err != nil {
×
114
                return err
×
115
        }
×
116

117
        return c.manager.Start(signals.SetupSignalHandler())
×
118
}
119

120
func (c *controller) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
×
121
        logger := slog.With("plan", req.Name)
×
122

×
123
        var plan api.KubeUpgradePlan
×
124
        err := c.Get(ctx, req.NamespacedName, &plan)
×
125
        if err != nil {
×
126
                logger.Error("Failed to get Plan", "err", err)
×
127
                return ctrl.Result{}, err
×
128
        }
×
129

130
        err = c.reconcile(ctx, &plan, logger)
×
131
        if err != nil {
×
132
                return ctrl.Result{}, err
×
133
        }
×
134

135
        err = c.Status().Update(ctx, &plan)
×
136
        if err != nil {
×
137
                logger.Error("Failed to update plan status", "err", err)
×
138
                return ctrl.Result{}, err
×
139
        }
×
140

141
        return ctrl.Result{
×
142
                Requeue:      plan.Status.Summary != api.PlanStatusComplete,
×
143
                RequeueAfter: time.Minute,
×
144
        }, nil
×
145
}
146

147
func (c *controller) reconcile(ctx context.Context, plan *api.KubeUpgradePlan, logger *slog.Logger) error {
16✔
148
        if plan.Status.Groups == nil {
26✔
149
                plan.Status.Groups = make(map[string]string, len(plan.Spec.Groups))
10✔
150
        }
10✔
151

152
        // Migration from v0.6.0: Remove the finalizer as it is not needed
153
        // TODO: Remove in future release
154
        if controllerutil.RemoveFinalizer(plan, constants.Finalizer) {
17✔
155
                logger.Debug("Removing finalizer from plan")
1✔
156
                err := c.Update(ctx, plan)
1✔
157
                if err != nil {
1✔
158
                        return fmt.Errorf("failed to remove finalizer from plan %s: %v", plan.Name, err)
×
159
                }
×
160
        }
161

162
        cmList := &corev1.ConfigMapList{}
16✔
163
        err := c.List(ctx, cmList, client.InNamespace(c.namespace), client.MatchingLabels{
16✔
164
                constants.LabelPlanName: plan.Name,
16✔
165
        })
16✔
166
        if err != nil {
16✔
167
                logger.Error("Failed to fetch upgraded ConfigMaps", "err", err)
×
168
                return err
×
169
        }
×
170

171
        dsList := &appv1.DaemonSetList{}
16✔
172
        err = c.List(ctx, dsList, client.InNamespace(c.namespace), client.MatchingLabels{
16✔
173
                constants.LabelPlanName: plan.Name,
16✔
174
        })
16✔
175
        if err != nil {
16✔
176
                logger.Error("Failed to fetch upgraded DaemonSets", "err", err)
×
177
                return err
×
178
        }
×
179

180
        daemons := make(map[string]*appv1.DaemonSet, len(plan.Spec.Groups))
16✔
181
        for i := range dsList.Items {
24✔
182
                daemon := &dsList.Items[i]
8✔
183
                group := daemon.Labels[constants.LabelNodeGroup]
8✔
184
                if _, ok := plan.Spec.Groups[group]; ok {
14✔
185
                        daemons[group] = daemon
6✔
186
                } else {
8✔
187
                        err = c.Delete(ctx, daemon)
2✔
188
                        if err != nil {
2✔
189
                                return fmt.Errorf("failed to delete DaemonSet %s: %v", daemon.Name, err)
×
190
                        }
×
191
                        logger.Info("Deleted obsolete DaemonSet", "name", daemon.Name)
2✔
192
                }
193
        }
194

195
        cms := make(map[string]*corev1.ConfigMap, len(plan.Spec.Groups))
16✔
196
        for i := range cmList.Items {
24✔
197
                cm := &cmList.Items[i]
8✔
198
                group := cm.Labels[constants.LabelNodeGroup]
8✔
199
                if _, ok := plan.Spec.Groups[group]; ok {
14✔
200
                        cms[group] = cm
6✔
201
                } else {
8✔
202
                        err = c.Delete(ctx, cm)
2✔
203
                        if err != nil {
2✔
204
                                return fmt.Errorf("failed to delete ConfigMap %s: %v", cm.Name, err)
×
205
                        }
×
206
                        logger.Info("Deleted obsolete ConfigMap", "name", cm.Name)
2✔
207
                }
208
        }
209

210
        nodesToUpdate := make(map[string][]corev1.Node, len(plan.Spec.Groups))
16✔
211
        newGroupStatus := make(map[string]string, len(plan.Spec.Groups))
16✔
212

16✔
213
        for name, cfg := range plan.Spec.Groups {
51✔
214
                logger := logger.With("group", name)
35✔
215

35✔
216
                err = c.reconcileUpgradedConfigMap(ctx, plan, logger, cms[name], name)
35✔
217
                if err != nil {
35✔
218
                        return fmt.Errorf("failed to reconcile ConfigMap for group %s: %v", name, err)
×
219
                }
×
220

221
                err = c.reconcileUpgradedDaemonSet(ctx, plan, logger, daemons[name], name, cfg)
35✔
222
                if err != nil {
35✔
223
                        return fmt.Errorf("failed to reconcile DaemonSet for group %s: %v", name, err)
×
224
                }
×
225

226
                nodeList := &corev1.NodeList{}
35✔
227
                err = c.List(ctx, nodeList, client.MatchingLabels(cfg.Labels))
35✔
228
                if err != nil {
35✔
229
                        logger.Error("Failed to get nodes for group", "err", err)
×
230
                        return err
×
231
                }
×
232

233
                status, update, nodes, err := c.reconcileNodes(plan.Spec.KubernetesVersion, plan.Spec.AllowDowngrade, nodeList.Items)
35✔
234
                if err != nil {
35✔
235
                        logger.Error("Failed to reconcile nodes for group", "err", err)
×
236
                        return err
×
237
                }
×
238

239
                newGroupStatus[name] = status
35✔
240

35✔
241
                if update {
60✔
242
                        nodesToUpdate[name] = nodes
25✔
243
                } else if plan.Status.Groups[name] != newGroupStatus[name] {
41✔
244
                        logger.Info("Group changed status", "status", newGroupStatus[name])
6✔
245
                }
6✔
246
        }
247

248
        for name, nodes := range nodesToUpdate {
41✔
249
                logger := logger.With("group", name)
25✔
250

25✔
251
                if groupWaitForDependency(plan.Spec.Groups[name].DependsOn, newGroupStatus) {
31✔
252
                        logger.Info("Group is waiting on dependencies")
6✔
253
                        newGroupStatus[name] = api.PlanStatusWaiting
6✔
254
                        continue
6✔
255
                } else if plan.Status.Groups[name] != newGroupStatus[name] {
38✔
256
                        logger.Info("Group changed status", "status", newGroupStatus[name])
19✔
257
                }
19✔
258

259
                for _, node := range nodes {
38✔
260
                        logger.Debug("Updating node annotations", "node", node.Name)
19✔
261
                        err = c.Update(ctx, &node)
19✔
262
                        if err != nil {
19✔
263
                                return fmt.Errorf("failed to update node %s: %v", node.GetName(), err)
×
264
                        }
×
265
                }
266
        }
267

268
        plan.Status.Groups = newGroupStatus
16✔
269
        plan.Status.Summary = createStatusSummary(plan.Status.Groups)
16✔
270

16✔
271
        return nil
16✔
272
}
273

274
func (c *controller) reconcileNodes(kubeVersion string, downgrade bool, nodes []corev1.Node) (string, bool, []corev1.Node, error) {
37✔
275
        if len(nodes) == 0 {
37✔
276
                return api.PlanStatusUnknown, false, nil, nil
×
277
        }
×
278

279
        completed := 0
37✔
280
        needUpdate := false
37✔
281
        errorNodes := make([]string, 0)
37✔
282

37✔
283
        for i := range nodes {
74✔
284
                if nodes[i].Annotations == nil {
61✔
285
                        nodes[i].Annotations = make(map[string]string)
24✔
286
                }
24✔
287

288
                if !downgrade && semver.Compare(kubeVersion, nodes[i].Status.NodeInfo.KubeletVersion) < 0 {
38✔
289
                        return api.PlanStatusError, false, nil, fmt.Errorf("node %s version %s is newer than %s, but downgrade is disabled", nodes[i].GetName(), nodes[i].Status.NodeInfo.KubeletVersion, kubeVersion)
1✔
290
                }
1✔
291

292
                if nodes[i].Annotations[constants.NodeKubernetesVersion] == kubeVersion {
46✔
293
                        switch nodes[i].Annotations[constants.NodeUpgradeStatus] {
10✔
294
                        case constants.NodeUpgradeStatusCompleted:
9✔
295
                                completed++
9✔
296
                        case constants.NodeUpgradeStatusError:
1✔
297
                                errorNodes = append(errorNodes, nodes[i].GetName())
1✔
298
                        }
299
                        continue
10✔
300
                }
301

302
                nodes[i].Annotations[constants.NodeKubernetesVersion] = kubeVersion
26✔
303
                nodes[i].Annotations[constants.NodeUpgradeStatus] = constants.NodeUpgradeStatusPending
26✔
304

26✔
305
                needUpdate = true
26✔
306
        }
307

308
        var status string
36✔
309
        if len(errorNodes) > 0 {
37✔
310
                status = fmt.Sprintf("%s: The nodes %v are reporting errors", api.PlanStatusError, errorNodes)
1✔
311
        } else if len(nodes) == completed {
45✔
312
                status = api.PlanStatusComplete
9✔
313
        } else {
35✔
314
                status = fmt.Sprintf("%s: %d/%d nodes upgraded", api.PlanStatusProgressing, completed, len(nodes))
26✔
315
        }
26✔
316
        return status, needUpdate, nodes, nil
36✔
317
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc