Skip to content
Open
Show file tree
Hide file tree
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
3 changes: 3 additions & 0 deletions pkg/dyntracker/dynamic_readiness_tracker.go
Original file line number Diff line number Diff line change
Expand Up @@ -1189,6 +1189,9 @@ func (t *DynamicReadinessTracker) handleDeploymentStatus(status *deployment.Depl
if status.IsFailed {
taskState.ResourceState(taskState.Name(), taskState.Namespace(), taskState.GroupVersionKind()).RWTransaction(func(rs *statestore.ResourceState) {
rs.AddError(errors.New(status.FailedReason), "", time.Now())
if status.FailureMode == deployment.FailureModeFatal {
rs.SetStatus(statestore.ResourceStatusFailed)
}
})
}
}
Expand Down
87 changes: 87 additions & 0 deletions pkg/dyntracker/dynamic_readiness_tracker_failure_mode_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
package dyntracker

import (
"testing"

"k8s.io/apimachinery/pkg/runtime/schema"

"github.com/werf/kubedog/pkg/dyntracker/statestore"
"github.com/werf/kubedog/pkg/tracker/deployment"
)

func TestHandleDeploymentStatus_KeepsCountedFailureWithinAllowance(t *testing.T) {
taskState := statestore.NewReadinessTaskState(
"demo",
"default",
schema.GroupVersionKind{Group: "apps", Version: "v1", Kind: "Deployment"},
statestore.ReadinessTaskStateOptions{TotalAllowFailuresCount: 1},
)
status := deployment.DeploymentStatus{
IsFailed: true,
FailedReason: "counted event failure",
FailureMode: deployment.FailureModeCounted,
}

(&DynamicReadinessTracker{}).handleDeploymentStatus(&status, taskState)

if got := taskState.Status(); got != statestore.ReadinessTaskStatusProgressing {
t.Fatalf("counted failure within allowance must keep tracking progressing, got %s", got)
}
}

func TestHandleDeploymentStatus_FailsOnFatalFailureInLegacyFailMode(t *testing.T) {
taskState := statestore.NewReadinessTaskState(
"demo",
"default",
schema.GroupVersionKind{Group: "apps", Version: "v1", Kind: "Deployment"},
statestore.ReadinessTaskStateOptions{
FailMode: statestore.LegacyHopeUntilEndOfDeployProcess,
TotalAllowFailuresCount: 1,
},
)
status := deployment.DeploymentStatus{
IsFailed: true,
FailedReason: `pods "demo-111" is forbidden: exceeded quota: compute-quota`,
FailureMode: deployment.FailureModeFatal,
}

(&DynamicReadinessTracker{}).handleDeploymentStatus(&status, taskState)

if got := taskState.Status(); got != statestore.ReadinessTaskStatusFailed {
t.Fatalf("fatal failure must fail tracking regardless of the allowance, got %s", got)
}
}

func TestHandleDeploymentStatus_KeepsFatalFailureIgnoredInIgnoreFailMode(t *testing.T) {
taskState := statestore.NewReadinessTaskState(
"demo",
"default",
schema.GroupVersionKind{Group: "apps", Version: "v1", Kind: "Deployment"},
statestore.ReadinessTaskStateOptions{
FailMode: statestore.IgnoreAndContinueDeployProcess,
TotalAllowFailuresCount: 0,
},
)
status := deployment.DeploymentStatus{
IsFailed: true,
FailedReason: `pods "demo-111" is forbidden: exceeded quota: compute-quota`,
FailureMode: deployment.FailureModeFatal,
}

(&DynamicReadinessTracker{}).handleDeploymentStatus(&status, taskState)

resourceState := taskState.ResourceState(taskState.Name(), taskState.Namespace(), taskState.GroupVersionKind())

var resourceStatus statestore.ResourceStatus
resourceState.RTransaction(func(rs *statestore.ResourceState) {
resourceStatus = rs.Status()
})

if resourceStatus != statestore.ResourceStatusFailed {
t.Fatalf("fatal failure must be recorded on the resource state, got %s", resourceStatus)
}

if got := taskState.Status(); got != statestore.ReadinessTaskStatusProgressing {
t.Fatalf("ignore fail mode must keep tracking progressing on a fatal failure, got %s", got)
}
}
180 changes: 180 additions & 0 deletions pkg/dyntracker/dynamic_readiness_tracker_replicafailure_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,180 @@
package dyntracker_test

import (
"context"
"errors"
"fmt"
"strings"
"testing"
"time"

appsv1 "k8s.io/api/apps/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"
dynamicfake "k8s.io/client-go/dynamic/fake"
k8sfake "k8s.io/client-go/kubernetes/fake"
k8sscheme "k8s.io/client-go/kubernetes/scheme"

"github.com/werf/kubedog/pkg/dyntracker"
"github.com/werf/kubedog/pkg/dyntracker/logstore"
"github.com/werf/kubedog/pkg/dyntracker/statestore"
"github.com/werf/kubedog/pkg/dyntracker/util"
"github.com/werf/kubedog/pkg/informer"
)

var replicaFailureDeploymentGVK = schema.GroupVersionKind{Group: "apps", Version: "v1", Kind: "Deployment"}

func TestDynamicReadinessTracker_SingleDurableReplicaFailure_TerminatesPromptly(t *testing.T) {
const boundedTimeout = 8 * time.Second

tests := []struct {
name string
allowFailures int
}{
{name: "allowance below the single failure", allowFailures: 0},
{name: "allowance equal to the single failure", allowFailures: 1},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
dynClient := dynamicfake.NewSimpleDynamicClient(k8sscheme.Scheme,
replicaFailureOneReplicaDeployment(), replicaFailureFailedCreateReplicaSet())
kubeClient := k8sfake.NewSimpleClientset()

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

watchErrCh := make(chan error, 100)
factory := informer.NewConcurrentInformerFactory(ctx.Done(), watchErrCh, dynClient, informer.ConcurrentInformerFactoryOptions{})

taskState := util.NewConcurrent(statestore.NewReadinessTaskState(
"demo", "default", replicaFailureDeploymentGVK,
statestore.ReadinessTaskStateOptions{TotalAllowFailuresCount: tt.allowFailures},
))
logStore := util.NewConcurrent(logstore.NewLogStore())

rt, err := dyntracker.NewDynamicReadinessTracker(
ctx, taskState, logStore, factory, kubeClient, dynClient, nil, replicaFailureStubRESTMapper{},
dyntracker.DynamicReadinessTrackerOptions{
Timeout: boundedTimeout,
CaseInsensitiveGVKMatching: true,
IgnoreLogs: true,
},
)
if err != nil {
t.Fatalf("NewDynamicReadinessTracker: %v", err)
}

start := time.Now()
trackErr := rt.Track(ctx)
elapsed := time.Since(start)

if trackErr == nil {
t.Fatalf("expected Track() to fail once the durable ReplicaFailure is observed, got nil error after %s", elapsed)
}
if !strings.Contains(trackErr.Error(), "readiness failed") {
t.Fatalf("expected a readiness-failed error, got: %v (after %s)", trackErr, elapsed)
}
if elapsed >= boundedTimeout/2 {
t.Fatalf("expected prompt termination well under the %s bailout bound, took %s: %v", boundedTimeout, elapsed, trackErr)
}
})
}
}

func replicaFailureOneReplicaDeployment() *appsv1.Deployment {
one := int32(1)
return &appsv1.Deployment{
TypeMeta: metav1.TypeMeta{Kind: "Deployment", APIVersion: "apps/v1"},
ObjectMeta: metav1.ObjectMeta{Name: "demo", Namespace: "default", UID: "dep-uid-1", Generation: 1},
Spec: appsv1.DeploymentSpec{
Replicas: &one,
Selector: &metav1.LabelSelector{MatchLabels: map[string]string{"app": "demo"}},
Template: corev1.PodTemplateSpec{
ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{"app": "demo"}},
Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "c", Image: "busybox:1"}}},
},
},
Status: appsv1.DeploymentStatus{ObservedGeneration: 1},
}
}

func replicaFailureFailedCreateReplicaSet() *appsv1.ReplicaSet {
one := int32(1)
isController := true
return &appsv1.ReplicaSet{
TypeMeta: metav1.TypeMeta{Kind: "ReplicaSet", APIVersion: "apps/v1"},
ObjectMeta: metav1.ObjectMeta{
Name: "demo-111",
Namespace: "default",
UID: "rs-uid-1",
Labels: map[string]string{"app": "demo"},
OwnerReferences: []metav1.OwnerReference{{
APIVersion: "apps/v1", Kind: "Deployment", Name: "demo", UID: "dep-uid-1", Controller: &isController,
}},
CreationTimestamp: metav1.Now(),
},
Spec: appsv1.ReplicaSetSpec{
Replicas: &one,
Selector: &metav1.LabelSelector{MatchLabels: map[string]string{"app": "demo"}},
Template: corev1.PodTemplateSpec{
ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{"app": "demo"}},
Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "c", Image: "busybox:1"}}},
},
},
Status: appsv1.ReplicaSetStatus{
Conditions: []appsv1.ReplicaSetCondition{{
Type: appsv1.ReplicaSetReplicaFailure,
Status: corev1.ConditionTrue,
Reason: "FailedCreate",
Message: `pods "demo-111" is forbidden: exceeded quota: compute-quota`,
LastTransitionTime: metav1.Now(),
}},
},
}
}

type replicaFailureStubRESTMapper struct{}

func (replicaFailureStubRESTMapper) KindFor(schema.GroupVersionResource) (schema.GroupVersionKind, error) {
return schema.GroupVersionKind{}, errors.New("not implemented")
}

func (replicaFailureStubRESTMapper) KindsFor(schema.GroupVersionResource) ([]schema.GroupVersionKind, error) {
return nil, errors.New("not implemented")
}

func (replicaFailureStubRESTMapper) ResourceFor(schema.GroupVersionResource) (schema.GroupVersionResource, error) {
return schema.GroupVersionResource{}, errors.New("not implemented")
}

func (replicaFailureStubRESTMapper) ResourcesFor(schema.GroupVersionResource) ([]schema.GroupVersionResource, error) {
return nil, errors.New("not implemented")
}

func (replicaFailureStubRESTMapper) RESTMapping(gk schema.GroupKind, _ ...string) (*meta.RESTMapping, error) {
if strings.EqualFold(gk.Group, "apps") && strings.EqualFold(gk.Kind, "deployment") {
return &meta.RESTMapping{
Resource: schema.GroupVersionResource{Group: "apps", Version: "v1", Resource: "deployments"},
GroupVersionKind: replicaFailureDeploymentGVK,
Scope: meta.RESTScopeNamespace,
}, nil
}
return nil, fmt.Errorf("replicaFailureStubRESTMapper: no mapping for %s", gk)
}

func (m replicaFailureStubRESTMapper) RESTMappings(gk schema.GroupKind, versions ...string) ([]*meta.RESTMapping, error) {
mapping, err := m.RESTMapping(gk, versions...)
if err != nil {
return nil, err
}
return []*meta.RESTMapping{mapping}, nil
}

func (replicaFailureStubRESTMapper) ResourceSingularizer(resource string) (string, error) {
return strings.TrimSuffix(resource, "s"), nil
}

func (replicaFailureStubRESTMapper) Reset() {}
42 changes: 23 additions & 19 deletions pkg/dyntracker/statestore/readiness_task_state.go
Original file line number Diff line number Diff line change
Expand Up @@ -190,29 +190,33 @@ func initReadinessTaskStateReadyConditions() []ReadinessTaskConditionFn {
func initReadinessTaskStateFailureConditions(failMode FailMode, totalAllowFailuresCount int) []ReadinessTaskConditionFn {
var failureConditions []ReadinessTaskConditionFn

failureConditions = append(failureConditions, func(taskState *ReadinessTaskState) bool {
var failed bool
lo.Must0(domigraph.BFS(taskState.resourceStatesTree, util.ResourceID(taskState.name, taskState.namespace, taskState.groupVersionKind), func(id string) bool {
state := lo.Must(taskState.resourceStatesTree.Vertex(id))
state.RTransaction(func(s *ResourceState) {
if s.Status() == ResourceStatusFailed {
failed = true
}
})
switch failMode {
case IgnoreAndContinueDeployProcess:
// This mode must do nothing when tracking of the resource fails, so neither a
// resource marked as failed nor an exceeded error budget may fail the task.
case FailWholeDeployProcessImmediately, LegacyHopeUntilEndOfDeployProcess:
// A resource marked as failed reports a durable failure it cannot recover from on
// its own, so it fails the task regardless of the allowed failures count.
failureConditions = append(failureConditions, func(taskState *ReadinessTaskState) bool {
var failed bool
lo.Must0(domigraph.BFS(taskState.resourceStatesTree, util.ResourceID(taskState.name, taskState.namespace, taskState.groupVersionKind), func(id string) bool {
state := lo.Must(taskState.resourceStatesTree.Vertex(id))
state.RTransaction(func(s *ResourceState) {
if s.Status() == ResourceStatusFailed {
failed = true
}
})

if failed {
return true
}
if failed {
return true
}

return false
}))
return false
}))

return failed
})
return failed
})

switch failMode {
case IgnoreAndContinueDeployProcess:
case FailWholeDeployProcessImmediately, LegacyHopeUntilEndOfDeployProcess:
failureConditions = append(failureConditions, func(taskState *ReadinessTaskState) bool {
var totalErrsCount int
lo.Must0(domigraph.BFS(taskState.resourceStatesTree, util.ResourceID(taskState.name, taskState.namespace, taskState.groupVersionKind), func(id string) bool {
Expand Down
23 changes: 23 additions & 0 deletions pkg/tracker/deployment/export_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
package deployment

import appsv1 "k8s.io/api/apps/v1"

func (d *Tracker) TestOnlyInjectReplicaSetAdded(rs *appsv1.ReplicaSet) {
d.replicaSetAdded <- rs
}

func (d *Tracker) TestOnlyInjectReplicaSetModified(rs *appsv1.ReplicaSet) {
d.replicaSetModified <- rs
}

func (d *Tracker) TestOnlyInjectReplicaSetDeleted(rs *appsv1.ReplicaSet) {
d.replicaSetDeleted <- rs
}

func (d *Tracker) TestOnlyReplicaSetDeletedChanCap() int {
return cap(d.replicaSetDeleted)
}

func (d *Tracker) TestOnlyInjectResourceFailure(reason string) {
d.resourceFailed <- reason
}
Loading
Loading