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
10 changes: 10 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,16 @@ jobs:
go install golang.org/x/tools/cmd/goimports@latest
files=$(goimports -l .); if [[ -n "$files" ]]; then echo "$files"; exit 1; fi
- run: go vet ./...
# go vet has no unchecked-type-assertion check. See .golangci.yml: an
# unchecked assertion in an informer handler is fatal, not recoverable.
# Must be built with Go >= the go directive in go.mod, otherwise it
# refuses to load the config ("the Go language version used to build
# golangci-lint is lower than the targeted Go version").
- name: golangci-lint (forcetypeassert)
uses: golangci/golangci-lint-action@v8
with:
version: v2.13.2
args: --timeout=8m
# /containers transitively imports github.com/NVIDIA/go-nvml. The
# bindings register NVML symbols (including some only present in
# very recent libnvidia-ml.so versions, e.g.
Expand Down
31 changes: 31 additions & 0 deletions .golangci.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
version: "2"

# Deliberately narrow. This config exists to gate one class of defect, not to
# impose a style regime on an existing codebase — a linter that fails on
# hundreds of pre-existing findings gets disabled instead of fixed.
#
# forcetypeassert: an unchecked type assertion panics on an unexpected type.
# In this agent that is not a recoverable error: informer handlers run on
# client-go's shared goroutine, where a panic reaches apimachinery's runtime
# handler and terminates the process. A customer cluster crashlooped node-agent
# 14 times in 12 hours on exactly one such assertion — a cache.DeletedFinalState
# Unknown tombstone asserted straight to *v1.Pod. go vet does not catch this.
linters:
default: none
enable:
- forcetypeassert

exclusions:
generated: lax
rules:
# Test code asserting on values it just constructed is not a production
# crash risk, and failing a test is already a visible outcome.
- path: _test\.go$
linters:
- forcetypeassert

issues:
# Every finding is a potential process-killing panic; do not collapse
# repeats of the same one.
max-issues-per-linter: 0
max-same-issues: 0
140 changes: 117 additions & 23 deletions common/ip_resolver.go
Original file line number Diff line number Diff line change
Expand Up @@ -376,13 +376,19 @@ func (resolver *K8sIPResolver) StartWatching() error {
func (resolve *K8sIPResolver) addReplicaSetHandlers(replicaSetInformer cache.SharedIndexInformer) {
replicaSetInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
rs := obj.(*appsv1.ReplicaSet)
rs, ok := objectAs[*appsv1.ReplicaSet](obj)
if !ok {
return
}
resolve.snapshot.ReplicaSets.Store(rs.UID, MinimalOwnerInfo{
OwnerReferences: rs.OwnerReferences,
})
},
UpdateFunc: func(oldObj, newObj interface{}) {
rs := newObj.(*appsv1.ReplicaSet)
rs, ok := objectAs[*appsv1.ReplicaSet](newObj)
if !ok {
return
}
resolve.snapshot.ReplicaSets.Store(rs.UID, MinimalOwnerInfo{
OwnerReferences: rs.OwnerReferences,
})
Expand Down Expand Up @@ -427,16 +433,43 @@ func deletedObject[T any](obj interface{}) (T, bool) {
return zero, false
}

// objectAs asserts an add/update informer payload to T without panicking.
//
// Add and update never carry a tombstone, so unlike deletedObject there is
// nothing to unwrap — but the assertion is still on the shared informer's
// goroutine, where a panic reaches apimachinery's runtime handler and kills the
// process. The informers here also register transform functions (stripPod,
// stripNode, stripService) that return the object unchanged when their own
// assertion fails, so an unexpected type can reach a handler rather than being
// filtered out.
//
// Skipping the event is the right failure mode: the resolver loses one object
// until the next resync, instead of taking the agent down.
func objectAs[T any](obj interface{}) (T, bool) {
if o, ok := obj.(T); ok {
return o, true
}
var zero T
klog.V(2).Infof("ignoring event with unexpected payload %T", obj)
return zero, false
}

func (resolve *K8sIPResolver) addDaemonSetHandlers(daemonSetInformer cache.SharedIndexInformer) {
daemonSetInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ds := obj.(*appsv1.DaemonSet)
ds, ok := objectAs[*appsv1.DaemonSet](obj)
if !ok {
return
}
resolve.snapshot.DaemonSets.Store(ds.UID, MinimalOwnerInfo{
OwnerReferences: ds.OwnerReferences,
})
},
UpdateFunc: func(oldObj, newObj interface{}) {
ds := newObj.(*appsv1.DaemonSet)
ds, ok := objectAs[*appsv1.DaemonSet](newObj)
if !ok {
return
}
resolve.snapshot.DaemonSets.Store(ds.UID, MinimalOwnerInfo{
OwnerReferences: ds.OwnerReferences,
})
Expand All @@ -454,13 +487,19 @@ func (resolve *K8sIPResolver) addDaemonSetHandlers(daemonSetInformer cache.Share
func (resolve *K8sIPResolver) addStatefulSetHandlers(statefulSetInformer cache.SharedIndexInformer) {
statefulSetInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ss := obj.(*appsv1.StatefulSet)
ss, ok := objectAs[*appsv1.StatefulSet](obj)
if !ok {
return
}
resolve.snapshot.StatefulSets.Store(ss.UID, MinimalOwnerInfo{
OwnerReferences: ss.OwnerReferences,
})
},
UpdateFunc: func(oldObj, newObj interface{}) {
ss := newObj.(*appsv1.StatefulSet)
ss, ok := objectAs[*appsv1.StatefulSet](newObj)
if !ok {
return
}
resolve.snapshot.StatefulSets.Store(ss.UID, MinimalOwnerInfo{
OwnerReferences: ss.OwnerReferences,
})
Expand All @@ -478,13 +517,19 @@ func (resolve *K8sIPResolver) addStatefulSetHandlers(statefulSetInformer cache.S
func (resolve *K8sIPResolver) addJobHandlers(jobInformer cache.SharedIndexInformer) {
jobInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
job := obj.(*batchv1.Job)
job, ok := objectAs[*batchv1.Job](obj)
if !ok {
return
}
resolve.snapshot.Jobs.Store(job.UID, MinimalOwnerInfo{
OwnerReferences: job.OwnerReferences,
})
},
UpdateFunc: func(oldObj, newObj interface{}) {
job := newObj.(*batchv1.Job)
job, ok := objectAs[*batchv1.Job](newObj)
if !ok {
return
}
resolve.snapshot.Jobs.Store(job.UID, MinimalOwnerInfo{
OwnerReferences: job.OwnerReferences,
})
Expand All @@ -502,13 +547,19 @@ func (resolve *K8sIPResolver) addJobHandlers(jobInformer cache.SharedIndexInform
func (resolve *K8sIPResolver) addCronJobHandlers(cronJobInformer cache.SharedIndexInformer) {
cronJobInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
cronJob := obj.(*batchv1.CronJob)
cronJob, ok := objectAs[*batchv1.CronJob](obj)
if !ok {
return
}
resolve.snapshot.CronJobs.Store(cronJob.UID, MinimalOwnerInfo{
OwnerReferences: cronJob.OwnerReferences,
})
},
UpdateFunc: func(oldObj, newObj interface{}) {
cronJob := newObj.(*batchv1.CronJob)
cronJob, ok := objectAs[*batchv1.CronJob](newObj)
if !ok {
return
}
resolve.snapshot.CronJobs.Store(cronJob.UID, MinimalOwnerInfo{
OwnerReferences: cronJob.OwnerReferences,
})
Expand All @@ -526,7 +577,10 @@ func (resolve *K8sIPResolver) addCronJobHandlers(cronJobInformer cache.SharedInd
func (resolve *K8sIPResolver) addServiceHandlers(serviceInformer cache.SharedIndexInformer) {
serviceInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
service := obj.(*v1.Service)
service, ok := objectAs[*v1.Service](obj)
if !ok {
return
}
minSvc := MinimalService{
Name: service.Name,
Namespace: service.Namespace,
Expand All @@ -542,8 +596,14 @@ func (resolve *K8sIPResolver) addServiceHandlers(serviceInformer cache.SharedInd
}
},
UpdateFunc: func(oldObj, newObj interface{}) {
oldService := oldObj.(*v1.Service)
service := newObj.(*v1.Service)
oldService, ok := objectAs[*v1.Service](oldObj)
if !ok {
return
}
service, ok := objectAs[*v1.Service](newObj)
if !ok {
return
}
minSvc := MinimalService{
Name: service.Name,
Namespace: service.Namespace,
Expand Down Expand Up @@ -586,13 +646,19 @@ func (resolve *K8sIPResolver) addServiceHandlers(serviceInformer cache.SharedInd
func (resolve *K8sIPResolver) addDeploymentHandlers(deploymentInformer cache.SharedIndexInformer) {
deploymentInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
deployment := obj.(*appsv1.Deployment)
deployment, ok := objectAs[*appsv1.Deployment](obj)
if !ok {
return
}
resolve.snapshot.Deployments.Store(deployment.UID, MinimalOwnerInfo{
OwnerReferences: deployment.OwnerReferences,
})
},
UpdateFunc: func(oldObj, newObj interface{}) {
deployment := newObj.(*appsv1.Deployment)
deployment, ok := objectAs[*appsv1.Deployment](newObj)
if !ok {
return
}
resolve.snapshot.Deployments.Store(deployment.UID, MinimalOwnerInfo{
OwnerReferences: deployment.OwnerReferences,
})
Expand All @@ -610,12 +676,21 @@ func (resolve *K8sIPResolver) addDeploymentHandlers(deploymentInformer cache.Sha
func (resolver *K8sIPResolver) addPodHandlers(podInformer cache.SharedIndexInformer) {
podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
pod := obj.(*v1.Pod)
pod, ok := objectAs[*v1.Pod](obj)
if !ok {
return
}
resolver.handlePodAdd(pod)
},
UpdateFunc: func(oldObj, newObj interface{}) {
oldPod := oldObj.(*v1.Pod)
newPod := newObj.(*v1.Pod)
oldPod, ok := objectAs[*v1.Pod](oldObj)
if !ok {
return
}
newPod, ok := objectAs[*v1.Pod](newObj)
if !ok {
return
}
// Clean old IPs that are no longer present
newIPs := make(map[string]bool, len(newPod.Status.PodIPs))
for _, ip := range newPod.Status.PodIPs {
Expand Down Expand Up @@ -686,14 +761,20 @@ func (resolver *K8sIPResolver) handlePodAdd(pod *v1.Pod) bool {
func (resolver *K8sIPResolver) addNodeHandlers(nodeInformer cache.SharedIndexInformer) {
nodeInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
node := obj.(*v1.Node)
node, ok := objectAs[*v1.Node](obj)
if !ok {
return
}
shouldReturn := resolver.handleNodeEvent(node)
if shouldReturn {
return
}
},
UpdateFunc: func(oldObj, newObj interface{}) {
node := newObj.(*v1.Node)
node, ok := objectAs[*v1.Node](newObj)
if !ok {
return
}
shouldReturn := resolver.handleNodeEvent(node)
if shouldReturn {
return
Expand Down Expand Up @@ -935,7 +1016,11 @@ func (resolver *K8sIPResolver) getControllerOfOwner(owner *metav1.OwnerReference
if !ok {
return nil, fmt.Errorf("%w: %s %s", errOwnerNotCached, owner.Kind, owner.UID)
}
info := val.(MinimalOwnerInfo)
info, ok := val.(MinimalOwnerInfo)
if !ok {
klog.V(5).Infof("type confusion in owner cache for %s %s", owner.Kind, owner.UID)
return nil, fmt.Errorf("%w: %s %s", errOwnerNotCached, owner.Kind, owner.UID)
}
return getControllerOwnerRef(info.OwnerReferences), nil
}

Expand Down Expand Up @@ -1113,8 +1198,17 @@ func (resolver *K8sIPResolver) resolvePodDescriptor(pod *MinimalPod) Workload {

func (resolver *K8sIPResolver) ResolvePodOwner(podName string, podNamespace string) Workload {
if uidVal, ok := resolver.snapshot.PodNameIndex.Load(podNamespace + "/" + podName); ok {
if podVal, ok := resolver.snapshot.Pods.Load(uidVal.(types.UID)); ok {
pod := podVal.(MinimalPod)
uid, ok := uidVal.(types.UID)
if !ok {
klog.V(5).Infof("type confusion in PodNameIndex for %s/%s", podNamespace, podName)
return Workload{}
}
if podVal, ok := resolver.snapshot.Pods.Load(uid); ok {
pod, ok := podVal.(MinimalPod)
if !ok {
klog.V(5).Infof("type confusion in Pods cache for %s/%s", podNamespace, podName)
return Workload{}
}
return resolver.resolvePodDescriptor(&pod)
}
Comment thread
mayankpande88 marked this conversation as resolved.
}
Expand Down
16 changes: 14 additions & 2 deletions containers/cilium.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,13 @@ func lookupCilium4(src, dst netaddr.IPPort) *netaddr.IPPort {
if err != nil || v == nil {
return nil
}
e := v.(*ctmap.CtEntry)
// Cilium BPF structs are decoded from bpffs, so a Cilium version whose
// layout differs from the one we build against can yield an unexpected
// type here. Degrade to "unresolved" rather than panicking the agent.
e, ok := v.(*ctmap.CtEntry)
if !ok {
return nil
}

backendKey := lbmap.NewBackend4KeyV3(loadbalancer.BackendID(e.BackendID))
b, err := backends4Map.Lookup(backendKey)
Expand Down Expand Up @@ -151,7 +157,13 @@ func lookupCilium6(src, dst netaddr.IPPort) *netaddr.IPPort {
if err != nil || v == nil {
return nil
}
e := v.(*ctmap.CtEntry)
// Cilium BPF structs are decoded from bpffs, so a Cilium version whose
// layout differs from the one we build against can yield an unexpected
// type here. Degrade to "unresolved" rather than panicking the agent.
e, ok := v.(*ctmap.CtEntry)
if !ok {
return nil
}
backendKey := lbmap.NewBackend6KeyV3(loadbalancer.BackendID(e.BackendID))
b, err := backends6Map.Lookup(backendKey)
if err != nil || b == nil {
Expand Down
14 changes: 13 additions & 1 deletion containers/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -416,7 +416,7 @@ func (c *Container) Collect(ch chan<- prometheus.Metric) {
for _, ctr := range p.parser.GetCounters() {
if ctr.Level == logparser.LevelCritical || ctr.Level == logparser.LevelError {
sample, _ := c.logSamples.LoadOrStore(ctr.Hash, common.TruncateUtf8(ctr.Sample, *flags.MaxLabelLength))
ch <- c.counter(metrics.LogMessages, float64(ctr.Messages), source, ctr.Level.String(), ctr.Hash, sample.(string))
ch <- c.counter(metrics.LogMessages, float64(ctr.Messages), source, ctr.Level.String(), ctr.Hash, sampleString(sample))
}
}
for _, sc := range p.parser.GetSensitiveCounters() {
Expand Down Expand Up @@ -2306,3 +2306,15 @@ func (c *Container) gauge(desc *prometheus.Desc, value float64, labelValues ...s
allLabels = append(allLabels, labelValues...)
return prometheus.MustNewConstMetric(desc, prometheus.GaugeValue, value, allLabels...)
}

// sampleValue renders a log sample stored in c.logSamples as a metric label.
//
// The map only ever holds strings, but this is a metric-collection path that
// runs on every scrape: a type confusion here should degrade the label, not
// take the agent down.
func sampleString(v interface{}) string {
if s, ok := v.(string); ok {
return s
}
return ""
}
8 changes: 7 additions & 1 deletion ebpftracer/tracer.go
Original file line number Diff line number Diff line change
Expand Up @@ -448,7 +448,13 @@ func getLostSamplesTracker(name string) *lostSamplesTracker {
if !ok {
tracker, _ = lostSamplesTrackers.LoadOrStore(name, &lostSamplesTracker{interval: 10})
}
return tracker.(*lostSamplesTracker)
t, ok := tracker.(*lostSamplesTracker)
if !ok {
// Only this function ever writes the map, so this is unreachable; return a
// throwaway rather than panicking in a metrics path.
return &lostSamplesTracker{interval: 10}
}
return t
}

// safeDuration converts a uint64 nanosecond value from eBPF to time.Duration.
Expand Down
6 changes: 5 additions & 1 deletion pinger/pinger.go
Original file line number Diff line number Diff line change
Expand Up @@ -214,7 +214,11 @@ func openConn() (*net.IPConn, error) {
if err != nil {
return nil, err
}
ipconn := conn.(*net.IPConn)
ipconn, ok := conn.(*net.IPConn)
if !ok {
conn.Close()
return nil, fmt.Errorf("unexpected connection type %T for ip4:icmp", conn)
}
f, err := ipconn.File()
if err != nil {
return nil, err
Expand Down
Loading