-
Notifications
You must be signed in to change notification settings - Fork 142
Expand file tree
/
Copy pathrollout_restart.go
More file actions
197 lines (179 loc) · 7.07 KB
/
Copy pathrollout_restart.go
File metadata and controls
197 lines (179 loc) · 7.07 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
// Copyright (c) HashiCorp, Inc.
// SPDX-License-Identifier: BUSL-1.1
package helpers
import (
"context"
"errors"
"fmt"
"strings"
"time"
argorolloutsv1alpha1 "github.com/argoproj/argo-rollouts/pkg/apis/rollouts/v1alpha1"
appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/client-go/tools/record"
ctrlclient "sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/log"
"github.com/hashicorp/vault-secrets-operator/api/v1beta1"
"github.com/hashicorp/vault-secrets-operator/consts"
)
// AnnotationRestartedAt is updated to trigger a rollout-restart
const AnnotationRestartedAt = "vso.secrets.hashicorp.com/restartedAt"
// strimziKafkaConnectGVK identifies the Strimzi KafkaConnect CRD targeted for
// rollout-restart. It is handled as an unstructured object so that VSO does not
// depend on the Strimzi Go API. Strimzi serves this resource under the
// kafka.strimzi.io/v1 API.
var strimziKafkaConnectGVK = schema.GroupVersionKind{
Group: "kafka.strimzi.io",
Version: "v1",
Kind: "KafkaConnect",
}
// HandleRolloutRestarts for all v1beta1.RolloutRestartTarget(s) configured for obj.
// Supported objs are: v1beta1.VaultDynamicSecret, v1beta1.VaultStaticSecret, v1beta1.VaultPKISecret
// Please note the following:
// - a rollout-restart will be triggered for each configured v1beta1.RolloutRestartTarget
// - the rollout-restart action has no support for roll-back
// - does not wait for the action to complete
//
// Returns all errors encountered.
func HandleRolloutRestarts(ctx context.Context, client ctrlclient.Client, obj ctrlclient.Object, recorder record.EventRecorder) error {
logger := log.FromContext(ctx)
var targets []v1beta1.RolloutRestartTarget
switch t := obj.(type) {
case *v1beta1.VaultDynamicSecret:
targets = t.Spec.RolloutRestartTargets
case *v1beta1.VaultStaticSecret:
targets = t.Spec.RolloutRestartTargets
case *v1beta1.VaultPKISecret:
targets = t.Spec.RolloutRestartTargets
case *v1beta1.HCPVaultSecretsApp:
targets = t.Spec.RolloutRestartTargets
default:
err := fmt.Errorf("unsupported Object type %T", t)
recorder.Eventf(obj, corev1.EventTypeWarning, consts.ReasonRolloutRestartUnsupported,
"Rollout restart impossible (please report this bug): err=%s", err)
return err
}
if len(targets) == 0 {
return nil
}
var errs error
for _, target := range targets {
if err := RolloutRestart(ctx, obj.GetNamespace(), target, client); err != nil {
errs = errors.Join(err)
recorder.Eventf(obj, corev1.EventTypeWarning, consts.ReasonRolloutRestartFailed,
"Rollout restart failed for target %#v: err=%s", target, err)
} else {
recorder.Eventf(obj, corev1.EventTypeNormal, consts.ReasonRolloutRestartTriggered,
"Rollout restart triggered for %v", target)
}
}
if errs != nil {
logger.Error(errs, "Rollout restart failed", "targets", targets)
} else {
logger.V(consts.LogLevelDebug).Info("Rollout restart succeeded", "total", len(targets))
}
return errs
}
// RolloutRestart patches the target in namespace for rollout-restart.
// Supported target Kinds are: DaemonSet, Deployment, StatefulSet, argo.Rollout, KafkaConnect
func RolloutRestart(ctx context.Context, namespace string, target v1beta1.RolloutRestartTarget, client ctrlclient.Client) error {
if namespace == "" {
return fmt.Errorf("namespace cannot be empty")
}
objectMeta := metav1.ObjectMeta{
Namespace: namespace,
Name: target.Name,
}
var obj ctrlclient.Object
switch target.Kind {
case "DaemonSet":
obj = &appsv1.DaemonSet{
ObjectMeta: objectMeta,
}
case "Deployment":
obj = &appsv1.Deployment{
ObjectMeta: objectMeta,
}
case "StatefulSet":
obj = &appsv1.StatefulSet{
ObjectMeta: objectMeta,
}
case "argo.Rollout":
obj = &argorolloutsv1alpha1.Rollout{
ObjectMeta: objectMeta,
}
case "KafkaConnect":
u := &unstructured.Unstructured{}
u.SetGroupVersionKind(strimziKafkaConnectGVK)
u.SetNamespace(namespace)
u.SetName(target.Name)
obj = u
default:
return fmt.Errorf("unsupported Kind %q for %T", target.Kind, target)
}
return patchForRolloutRestart(ctx, obj, client)
}
func patchForRolloutRestart(ctx context.Context, obj ctrlclient.Object, client ctrlclient.Client) error {
objKey := ctrlclient.ObjectKeyFromObject(obj)
if err := client.Get(ctx, objKey, obj); err != nil {
return fmt.Errorf("failed to Get object for objKey %s, err=%w", objKey, err)
}
switch t := obj.(type) {
case *appsv1.Deployment:
if t.Spec.Paused {
return fmt.Errorf("deployment %s is paused, cannot restart it", obj)
}
patch := ctrlclient.StrategicMergeFrom(t.DeepCopy())
if t.Spec.Template.ObjectMeta.Annotations == nil {
t.Spec.Template.ObjectMeta.Annotations = make(map[string]string)
}
t.Spec.Template.ObjectMeta.Annotations[AnnotationRestartedAt] = time.Now().Format(time.RFC3339)
return client.Patch(ctx, t, patch)
case *appsv1.StatefulSet:
patch := ctrlclient.StrategicMergeFrom(t.DeepCopy())
if t.Spec.Template.ObjectMeta.Annotations == nil {
t.Spec.Template.ObjectMeta.Annotations = make(map[string]string)
}
t.Spec.Template.ObjectMeta.Annotations[AnnotationRestartedAt] = time.Now().Format(time.RFC3339)
return client.Patch(ctx, t, patch)
case *appsv1.DaemonSet:
patch := ctrlclient.StrategicMergeFrom(t.DeepCopy())
if t.Spec.Template.ObjectMeta.Annotations == nil {
t.Spec.Template.ObjectMeta.Annotations = make(map[string]string)
}
t.Spec.Template.ObjectMeta.Annotations[AnnotationRestartedAt] = time.Now().Format(time.RFC3339)
return client.Patch(ctx, t, patch)
case *argorolloutsv1alpha1.Rollout:
// use MergeFrom() since it supports CRDs whereas StrategicMergeFrom() does not.
patch := ctrlclient.MergeFrom(t.DeepCopy())
t.Spec.RestartAt = &metav1.Time{Time: time.Now()}
return client.Patch(ctx, t, patch)
case *unstructured.Unstructured:
// Strimzi CRDs (e.g. KafkaConnect) are handled as unstructured objects so
// that VSO does not depend on the Strimzi Go API. Strimzi triggers a
// rolling update of the managed pods when the annotations under
// spec.template.pod.metadata.annotations change. Use MergeFrom() since it
// supports CRDs whereas StrategicMergeFrom() does not.
annotationsPath := []string{"spec", "template", "pod", "metadata", "annotations"}
patch := ctrlclient.MergeFrom(t.DeepCopy())
annotations, _, err := unstructured.NestedStringMap(t.Object, annotationsPath...)
if err != nil {
return fmt.Errorf("failed to read %s for %s, err=%w",
strings.Join(annotationsPath, "."), ctrlclient.ObjectKeyFromObject(t), err)
}
if annotations == nil {
annotations = make(map[string]string)
}
annotations[AnnotationRestartedAt] = time.Now().Format(time.RFC3339)
if err := unstructured.SetNestedStringMap(t.Object, annotations, annotationsPath...); err != nil {
return fmt.Errorf("failed to set %s for %s, err=%w",
strings.Join(annotationsPath, "."), ctrlclient.ObjectKeyFromObject(t), err)
}
return client.Patch(ctx, t, patch)
default:
return fmt.Errorf("unsupported type %T for rollout-restart patching", t)
}
}