Skip to content

Commit a414e31

Browse files
authored
fix: keep recoverable ops failures in progress (#10299)
1 parent 37a10b7 commit a414e31

6 files changed

Lines changed: 181 additions & 18 deletions

File tree

pkg/operations/ops_comp_helper.go

Lines changed: 21 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -135,13 +135,14 @@ func (c componentOpsHelper) cancelComponentOps(ctx context.Context,
135135
return cli.Update(ctx, opsRes.Cluster)
136136
}
137137

138-
func (c componentOpsHelper) existFailure(ops *opsv1alpha1.OpsRequest, componentName string) bool {
139-
for _, v := range ops.Status.Components[componentName].ProgressDetails {
138+
func componentStatusFailureCount(compStatus opsv1alpha1.OpsRequestComponentStatus) int32 {
139+
var count int32
140+
for _, v := range compStatus.ProgressDetails {
140141
if v.Status == opsv1alpha1.FailedProgressStatus {
141-
return true
142+
count++
142143
}
143144
}
144-
return false
145+
return count
145146
}
146147

147148
func (c componentOpsHelper) getComponentOps(componentName string) (ComponentOpsInterface, bool) {
@@ -263,30 +264,34 @@ func (c componentOpsHelper) reconcileActionWithComponentOps(reqCtx intctrlutil.R
263264
existFailure := false
264265
for i := range progressResources {
265266
pgResource := progressResources[i]
267+
var componentPhase appsv1.ComponentPhase
268+
if pgResource.shards == nil {
269+
componentPhase = opsRes.Cluster.Status.Components[pgResource.compOps.GetComponentName()].Phase
270+
} else {
271+
componentPhase = opsRes.Cluster.Status.Shardings[pgResource.compOps.GetComponentName()].Phase
272+
}
273+
pgResource.componentPhase = componentPhase
266274
opsCompStatus := opsRequest.Status.Components[pgResource.compOps.GetComponentName()]
267275
expectCount, completedCount, err := handleStatusProgress(reqCtx, cli, opsRes, &pgResource, &opsCompStatus)
268276
if err != nil {
269277
return opsRequestPhase, 0, err
270278
}
271-
expectProgressCount += expectCount
272-
completedProgressCount += completedCount
273-
if c.existFailure(opsRes.OpsRequest, pgResource.compOps.GetComponentName()) {
279+
componentFailureCount := componentStatusFailureCount(opsCompStatus)
280+
componentHasFailure := componentFailureCount > 0
281+
if componentHasFailure {
274282
existFailure = true
275283
}
276-
var componentPhase appsv1.ComponentPhase
277-
if pgResource.shards == nil {
278-
componentPhase = opsRes.Cluster.Status.Components[pgResource.compOps.GetComponentName()].Phase
279-
} else {
280-
componentPhase = opsRes.Cluster.Status.Shardings[pgResource.compOps.GetComponentName()].Phase
281-
}
284+
expectProgressCount += expectCount
285+
completedProgressCount += completedCount
282286
// conditions whether ops is running:
283287
// 1. completedProgressCount is not equal to expectProgressCount.
284288
// 2. the component phase is not a terminal phase or no completed progress if the ops
285289
// needs to wait for the component phase to reach a terminal state.
286-
if expectCount != completedCount {
290+
switch {
291+
case expectCount != completedCount:
287292
opsIsCompleted = false
288-
} else if !pgResource.noWaitComponentCompleted &&
289-
(!slices.Contains(componentTerminalPhases(), componentPhase) || noAnyProgressCompleted(pgResource.clusterComponent.Replicas, completedCount)) {
293+
case !pgResource.noWaitComponentCompleted &&
294+
(!slices.Contains(componentTerminalPhases(), componentPhase) || noAnyProgressCompleted(pgResource.clusterComponent.Replicas, completedCount)):
290295
opsIsCompleted = false
291296
}
292297
opsCompStatus.Phase = componentPhase

pkg/operations/ops_manager.go

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -274,10 +274,14 @@ func (opsMgr *OpsManager) checkAndHandleOpsTimeout(reqCtx intctrlutil.RequestCtx
274274
return 0, PatchOpsStatus(reqCtx.Ctx, cli, opsRes, opsv1alpha1.OpsAbortedPhase,
275275
opsv1alpha1.NewAbortedCondition("Aborted due to exceeding the specified timeout period (timeoutSeconds)"))
276276
}
277+
timeoutRequeueAfter := time.Until(timeoutPoint)
277278
if requeueAfter != 0 {
279+
if timeoutRequeueAfter < requeueAfter {
280+
return timeoutRequeueAfter, nil
281+
}
278282
return requeueAfter, nil
279283
}
280-
return time.Until(timeoutPoint), nil
284+
return timeoutRequeueAfter, nil
281285
}
282286

283287
func GetOpsManager() *OpsManager {

pkg/operations/ops_progress_util.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -320,7 +320,8 @@ func handleFailedOrProcessingProgressDetail(opsRes *OpsResource,
320320
progressDetail opsv1alpha1.ProgressStatusDetail,
321321
instance Instance) (completedCount int32) {
322322
componentName := pgRes.clusterComponent.Name
323-
if instance.IsFailedAndTimedOut() {
323+
if pgRes.componentPhase == appsv1.FailedComponentPhase ||
324+
(instance.IsFailedAndTimedOut() && !pgRes.deferInstanceFailureToWorkloadPhase) {
324325
podMessage := getFailedPodMessage(opsRes.Cluster, componentName, instance.GetName())
325326
message := getProgressFailedMessage(pgRes.opsMessageKey, progressDetail.ObjectKey, componentName, podMessage)
326327
progressDetail.SetStatusAndMessage(opsv1alpha1.FailedProgressStatus, message)

pkg/operations/ops_util_test.go

Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -151,6 +151,154 @@ var _ = Describe("OpsUtil functions", func() {
151151
Expect(opsPhase).Should(Equal(opsv1alpha1.OpsFailedPhase))
152152
})
153153

154+
It("keeps restart ops running when a failed progress is not backed by a failed component", func() {
155+
By("init operations resources ")
156+
opsRes, _, _ := initOperationsResources(compDefName, clusterName)
157+
testapps.MockInstanceSetComponent(&testCtx, clusterName, defaultCompName)
158+
pods := testapps.MockInstanceSetPods(&testCtx, nil, opsRes.Cluster, defaultCompName)
159+
time.Sleep(time.Second)
160+
161+
ops := testops.NewOpsRequestObj("restart-ops-"+randomStr, testCtx.DefaultNamespace,
162+
clusterName, opsv1alpha1.RestartType)
163+
ops.Spec.RestartList = []opsv1alpha1.ComponentOps{{ComponentName: defaultCompName}}
164+
opsRes.OpsRequest = testops.CreateOpsRequest(ctx, testCtx, ops)
165+
opsRes.OpsRequest.Status.Phase = opsv1alpha1.OpsRunningPhase
166+
opsRes.OpsRequest.Status.StartTimestamp = metav1.NewTime(time.Now().Add(-time.Second))
167+
runtimes, err := buildOpsRuntimes(ctx, k8sClient, opsRes)
168+
Expect(err).Should(BeNil())
169+
opsRes.Runtimes = runtimes
170+
171+
handleRestartProgress := func(reqCtx intctrlutil.RequestCtx,
172+
cli client.Client,
173+
opsRes *OpsResource,
174+
pgRes *progressResource,
175+
compStatus *opsv1alpha1.OpsRequestComponentStatus) (expectProgressCount int32, completedCount int32, err error) {
176+
pgRes.deferInstanceFailureToWorkloadPhase = true
177+
return handleComponentStatusProgress(reqCtx, cli, opsRes, pgRes, compStatus,
178+
func(ops *opsv1alpha1.OpsRequest, instance Instance, pgRes *progressResource) bool {
179+
creationTimestamp := instance.GetCreationTimestamp()
180+
return !creationTimestamp.Before(&ops.Status.StartTimestamp)
181+
})
182+
}
183+
184+
recreatePod := func(pod *corev1.Pod) *corev1.Pod {
185+
testk8s.MockPodIsTerminating(ctx, testCtx, pod)
186+
testk8s.RemovePodFinalizer(ctx, testCtx, pod)
187+
return testapps.MockInstanceSetPod(&testCtx, nil, clusterName, defaultCompName, pod.Name, "follower")
188+
}
189+
for i := range pods {
190+
pods[i] = recreatePod(pods[i])
191+
}
192+
testk8s.MockPodIsFailed(ctx, testCtx, pods[2])
193+
194+
reqCtx := intctrlutil.RequestCtx{Ctx: ctx}
195+
compOpsHelper := newComponentOpsHelper(opsRes.OpsRequest.Spec.RestartList)
196+
opsPhase, requeueAfter, err := compOpsHelper.reconcileActionWithComponentOps(reqCtx, k8sClient, opsRes,
197+
"test", handleRestartProgress)
198+
Expect(err).Should(BeNil())
199+
Expect(opsPhase).Should(Equal(opsv1alpha1.OpsRunningPhase))
200+
Expect(requeueAfter).Should(BeZero())
201+
Expect(opsRes.OpsRequest.Status.Progress).Should(Equal("2/3"))
202+
progressDetail := findStatusProgressDetail(opsRes.OpsRequest.Status.Components[defaultCompName].ProgressDetails,
203+
getProgressObjectKey(constant.PodKind, pods[2].Name))
204+
Expect(progressDetail).ShouldNot(BeNil())
205+
Expect(progressDetail.Status).Should(Equal(opsv1alpha1.ProcessingProgressStatus))
206+
207+
By("mock the failed instance recovers while the component remains Running")
208+
recoveredPod := pods[2]
209+
patch := client.MergeFrom(recoveredPod.DeepCopy())
210+
recoveredPod.Status.Conditions = []corev1.PodCondition{
211+
{
212+
Type: corev1.PodReady,
213+
Status: corev1.ConditionTrue,
214+
},
215+
{
216+
Type: corev1.ContainersReady,
217+
Status: corev1.ConditionTrue,
218+
},
219+
}
220+
recoveredPod.Status.ContainerStatuses = []corev1.ContainerStatus{
221+
{
222+
Name: recoveredPod.Spec.Containers[0].Name,
223+
State: corev1.ContainerState{
224+
Running: &corev1.ContainerStateRunning{},
225+
},
226+
},
227+
}
228+
Expect(k8sClient.Status().Patch(ctx, recoveredPod, patch)).Should(Succeed())
229+
230+
opsPhase, _, err = compOpsHelper.reconcileActionWithComponentOps(reqCtx, k8sClient, opsRes,
231+
"test", handleRestartProgress)
232+
Expect(err).Should(BeNil())
233+
Expect(opsPhase).Should(Equal(opsv1alpha1.OpsSucceedPhase))
234+
Expect(opsRes.OpsRequest.Status.Progress).Should(Equal("3/3"))
235+
})
236+
237+
It("fails restart ops when component reaches terminal failed phase", func() {
238+
By("init operations resources ")
239+
opsRes, _, _ := initOperationsResources(compDefName, clusterName)
240+
testapps.MockInstanceSetComponent(&testCtx, clusterName, defaultCompName)
241+
pods := testapps.MockInstanceSetPods(&testCtx, nil, opsRes.Cluster, defaultCompName)
242+
time.Sleep(time.Second)
243+
244+
ops := testops.NewOpsRequestObj("restart-ops-"+randomStr, testCtx.DefaultNamespace,
245+
clusterName, opsv1alpha1.RestartType)
246+
ops.Spec.RestartList = []opsv1alpha1.ComponentOps{{ComponentName: defaultCompName}}
247+
opsRes.OpsRequest = testops.CreateOpsRequest(ctx, testCtx, ops)
248+
opsRes.OpsRequest.Status.Phase = opsv1alpha1.OpsRunningPhase
249+
opsRes.OpsRequest.Status.StartTimestamp = metav1.NewTime(time.Now().Add(-time.Second))
250+
runtimes, err := buildOpsRuntimes(ctx, k8sClient, opsRes)
251+
Expect(err).Should(BeNil())
252+
opsRes.Runtimes = runtimes
253+
254+
handleRestartProgress := func(reqCtx intctrlutil.RequestCtx,
255+
cli client.Client,
256+
opsRes *OpsResource,
257+
pgRes *progressResource,
258+
compStatus *opsv1alpha1.OpsRequestComponentStatus) (expectProgressCount int32, completedCount int32, err error) {
259+
pgRes.deferInstanceFailureToWorkloadPhase = true
260+
return handleComponentStatusProgress(reqCtx, cli, opsRes, pgRes, compStatus,
261+
func(ops *opsv1alpha1.OpsRequest, instance Instance, pgRes *progressResource) bool {
262+
creationTimestamp := instance.GetCreationTimestamp()
263+
return !creationTimestamp.Before(&ops.Status.StartTimestamp)
264+
})
265+
}
266+
267+
recreatePod := func(pod *corev1.Pod) *corev1.Pod {
268+
testk8s.MockPodIsTerminating(ctx, testCtx, pod)
269+
testk8s.RemovePodFinalizer(ctx, testCtx, pod)
270+
return testapps.MockInstanceSetPod(&testCtx, nil, clusterName, defaultCompName, pod.Name, "follower")
271+
}
272+
for i := range pods {
273+
pods[i] = recreatePod(pods[i])
274+
}
275+
testk8s.MockPodIsFailed(ctx, testCtx, pods[2])
276+
277+
reqCtx := intctrlutil.RequestCtx{Ctx: ctx}
278+
compOpsHelper := newComponentOpsHelper(opsRes.OpsRequest.Spec.RestartList)
279+
opsPhase, requeueAfter, err := compOpsHelper.reconcileActionWithComponentOps(reqCtx, k8sClient, opsRes,
280+
"test", handleRestartProgress)
281+
Expect(err).Should(BeNil())
282+
Expect(opsPhase).Should(Equal(opsv1alpha1.OpsRunningPhase))
283+
Expect(requeueAfter).Should(BeZero())
284+
285+
By("mock component reaches terminal Failed phase")
286+
clusterComp := opsRes.Cluster.Status.Components[defaultCompName]
287+
clusterComp.Phase = appsv1.FailedComponentPhase
288+
opsRes.Cluster.Status.SetComponentStatus(defaultCompName, clusterComp)
289+
290+
opsPhase, requeueAfter, err = compOpsHelper.reconcileActionWithComponentOps(reqCtx, k8sClient, opsRes,
291+
"test", handleRestartProgress)
292+
Expect(err).Should(BeNil())
293+
Expect(requeueAfter).Should(BeZero())
294+
Expect(opsPhase).Should(Equal(opsv1alpha1.OpsFailedPhase))
295+
Expect(opsRes.OpsRequest.Status.Progress).Should(Equal("3/3"))
296+
progressDetail := findStatusProgressDetail(opsRes.OpsRequest.Status.Components[defaultCompName].ProgressDetails,
297+
getProgressObjectKey(constant.PodKind, pods[2].Name))
298+
Expect(progressDetail).ShouldNot(BeNil())
299+
Expect(progressDetail.Status).Should(Equal(opsv1alpha1.FailedProgressStatus))
300+
})
301+
154302
It("Test opsRequest with disable ha", func() {
155303
By("init operations resources ")
156304
opsRes, _, _ := initOperationsResources(compDefName, clusterName)

pkg/operations/restart.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,7 @@ func (r restartOpsHandler) ReconcileAction(reqCtx intctrlutil.RequestCtx, cli cl
8989
opsRes *OpsResource,
9090
pgRes *progressResource,
9191
compStatus *opsv1alpha1.OpsRequestComponentStatus) (expectProgressCount int32, completedCount int32, err error) {
92+
pgRes.deferInstanceFailureToWorkloadPhase = true
9293
return handleComponentStatusProgress(reqCtx, cli, opsRes, pgRes, compStatus, r.podApplyCompOps)
9394
}
9495
return r.compOpsHelper.reconcileActionWithComponentOps(reqCtx, cli, opsRes,

pkg/operations/type.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,10 @@ type progressResource struct {
111111
// checks if it needs to wait the component to complete.
112112
// if only updates a part of pods, set it to false.
113113
noWaitComponentCompleted bool
114+
// lets ops types such as restart defer pod-level failure signals until the
115+
// workload/component reaches a terminal failure state.
116+
deferInstanceFailureToWorkloadPhase bool
117+
componentPhase appsv1.ComponentPhase
114118
}
115119

116120
// OpsRuntime abstracts the standard ops paths that only need workload/member views

0 commit comments

Comments
 (0)