Skip to content
Merged
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
75 changes: 66 additions & 9 deletions common/ip_resolver.go
Original file line number Diff line number Diff line change
Expand Up @@ -388,12 +388,45 @@ func (resolve *K8sIPResolver) addReplicaSetHandlers(replicaSetInformer cache.Sha
})
},
DeleteFunc: func(obj interface{}) {
rs := obj.(*appsv1.ReplicaSet)
rs, ok := deletedObject[*appsv1.ReplicaSet](obj)
if !ok {
return
}
resolve.snapshot.ReplicaSets.Delete(rs.UID)
},
})
}

// deletedObject extracts the deleted object from an informer DeleteFunc payload.
//
// client-go does not always deliver the object itself on a delete. When the
// watch is interrupted and the final delete event is missed, the informer
// resyncs and delivers a cache.DeletedFinalStateUnknown tombstone wrapping the
// last known state. Asserting the payload directly panics on that tombstone,
// and because these handlers run on the shared informer's goroutine the panic
// reaches k8s.io/apimachinery's runtime handler, which is fatal — the whole
// agent exits and the pod restarts.
//
// Reported from a customer cluster where node-agent crashlooped 15 times:
//
// panic: interface conversion: interface {} is cache.DeletedFinalStateUnknown, not *v1.Pod
//
// A tombstone whose payload is the wrong type, or nil, is not recoverable — the
// caller skips that delete rather than dying.
func deletedObject[T any](obj interface{}) (T, bool) {
if o, ok := obj.(T); ok {
return o, true
}
if tombstone, ok := obj.(cache.DeletedFinalStateUnknown); ok {
if o, ok := tombstone.Obj.(T); ok {
return o, true
}
}
var zero T
klog.V(2).Infof("ignoring delete event with unexpected payload %T", obj)
return zero, false
}

func (resolve *K8sIPResolver) addDaemonSetHandlers(daemonSetInformer cache.SharedIndexInformer) {
daemonSetInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
Expand All @@ -409,7 +442,10 @@ func (resolve *K8sIPResolver) addDaemonSetHandlers(daemonSetInformer cache.Share
})
},
DeleteFunc: func(obj interface{}) {
ds := obj.(*appsv1.DaemonSet)
ds, ok := deletedObject[*appsv1.DaemonSet](obj)
if !ok {
return
}
resolve.snapshot.DaemonSets.Delete(ds.UID)
},
})
Expand All @@ -430,7 +466,10 @@ func (resolve *K8sIPResolver) addStatefulSetHandlers(statefulSetInformer cache.S
})
},
DeleteFunc: func(obj interface{}) {
ss := obj.(*appsv1.StatefulSet)
ss, ok := deletedObject[*appsv1.StatefulSet](obj)
if !ok {
return
}
resolve.snapshot.StatefulSets.Delete(ss.UID)
},
})
Expand All @@ -451,7 +490,10 @@ func (resolve *K8sIPResolver) addJobHandlers(jobInformer cache.SharedIndexInform
})
},
DeleteFunc: func(obj interface{}) {
job := obj.(*batchv1.Job)
job, ok := deletedObject[*batchv1.Job](obj)
if !ok {
return
}
resolve.snapshot.Jobs.Delete(job.UID)
},
})
Expand All @@ -472,7 +514,10 @@ func (resolve *K8sIPResolver) addCronJobHandlers(cronJobInformer cache.SharedInd
})
},
DeleteFunc: func(obj interface{}) {
cronJob := obj.(*batchv1.CronJob)
cronJob, ok := deletedObject[*batchv1.CronJob](obj)
if !ok {
return
}
resolve.snapshot.CronJobs.Delete(cronJob.UID)
},
})
Expand Down Expand Up @@ -524,7 +569,10 @@ func (resolve *K8sIPResolver) addServiceHandlers(serviceInformer cache.SharedInd
}
},
DeleteFunc: func(obj interface{}) {
service := obj.(*v1.Service)
service, ok := deletedObject[*v1.Service](obj)
if !ok {
return
}
resolve.snapshot.Services.Delete(service.UID)
for _, clusterIp := range service.Spec.ClusterIPs {
if clusterIp != "None" {
Expand All @@ -550,7 +598,10 @@ func (resolve *K8sIPResolver) addDeploymentHandlers(deploymentInformer cache.Sha
})
},
DeleteFunc: func(obj interface{}) {
deployment := obj.(*appsv1.Deployment)
deployment, ok := deletedObject[*appsv1.Deployment](obj)
if !ok {
return
}
resolve.snapshot.Deployments.Delete(deployment.UID)
},
})
Expand Down Expand Up @@ -578,7 +629,10 @@ func (resolver *K8sIPResolver) addPodHandlers(podInformer cache.SharedIndexInfor
resolver.handlePodAdd(newPod)
},
DeleteFunc: func(obj interface{}) {
pod := obj.(*v1.Pod)
pod, ok := deletedObject[*v1.Pod](obj)
if !ok {
return
}
resolver.snapshot.Pods.Delete(pod.UID)
resolver.snapshot.PodDescriptors.Delete(pod.UID)
resolver.snapshot.PodNameIndex.Delete(pod.Namespace + "/" + pod.Name)
Expand Down Expand Up @@ -646,7 +700,10 @@ func (resolver *K8sIPResolver) addNodeHandlers(nodeInformer cache.SharedIndexInf
}
},
DeleteFunc: func(obj interface{}) {
node := obj.(*v1.Node)
node, ok := deletedObject[*v1.Node](obj)
if !ok {
return
}
resolver.snapshot.Nodes.Delete(node.UID)
resolver.instanceMetaMap.Delete(node.Name)
},
Expand Down
77 changes: 77 additions & 0 deletions common/ip_resolver_delete_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
package common

import (
"testing"

v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/tools/cache"
)

// client-go delivers a cache.DeletedFinalStateUnknown tombstone instead of the
// object when a watch is interrupted and the final delete is missed. The
// informer DeleteFuncs used to assert the payload directly, so a tombstone
// panicked on the shared informer goroutine — fatal to the process. A customer
// cluster crashlooped node-agent 15 times on exactly this:
//
// panic: interface conversion: interface {} is cache.DeletedFinalStateUnknown, not *v1.Pod
func TestDeletedObjectUnwrapsTombstone(t *testing.T) {
pod := &v1.Pod{ObjectMeta: metav1.ObjectMeta{Name: "p", Namespace: "ns", UID: "uid-1"}}

t.Run("plain object", func(t *testing.T) {
got, ok := deletedObject[*v1.Pod](pod)
if !ok || got != pod {
t.Fatalf("plain object not returned: got=%v ok=%v", got, ok)
}
})

t.Run("tombstone", func(t *testing.T) {
got, ok := deletedObject[*v1.Pod](cache.DeletedFinalStateUnknown{Key: "ns/p", Obj: pod})
if !ok {
t.Fatal("tombstone not unwrapped — this is the crash-loop bug")
}
if got != pod {
t.Fatalf("wrong object from tombstone: %v", got)
}
})
}

// A payload we cannot make sense of must be skipped, never panic: a tombstone
// wrapping the wrong type or nil, or an unrelated object.
func TestDeletedObjectRejectsBadPayloads(t *testing.T) {
cases := []struct {
name string
obj interface{}
}{
{"nil", nil},
{"wrong type", &v1.Service{}},
{"tombstone wrapping wrong type", cache.DeletedFinalStateUnknown{Key: "k", Obj: &v1.Service{}}},
{"tombstone wrapping nil", cache.DeletedFinalStateUnknown{Key: "k", Obj: nil}},
{"unrelated value", 42},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
defer func() {
if r := recover(); r != nil {
t.Fatalf("panicked on %s: %v", tc.name, r)
}
}()
if got, ok := deletedObject[*v1.Pod](tc.obj); ok {
t.Fatalf("accepted bad payload %s: %v", tc.name, got)
}
})
}
}

// The same tombstone path must hold for every object type with a DeleteFunc,
// not just pods — Service and Node take the same shape.
func TestDeletedObjectOtherTypes(t *testing.T) {
svc := &v1.Service{ObjectMeta: metav1.ObjectMeta{Name: "s", UID: "uid-s"}}
if got, ok := deletedObject[*v1.Service](cache.DeletedFinalStateUnknown{Obj: svc}); !ok || got != svc {
t.Error("service tombstone not unwrapped")
}
node := &v1.Node{ObjectMeta: metav1.ObjectMeta{Name: "n", UID: "uid-n"}}
if got, ok := deletedObject[*v1.Node](cache.DeletedFinalStateUnknown{Obj: node}); !ok || got != node {
t.Error("node tombstone not unwrapped")
}
}
Loading