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
2 changes: 2 additions & 0 deletions docs/mcp.md
Original file line number Diff line number Diff line change
Expand Up @@ -575,3 +575,5 @@ other writers, and later template/configuration edits can explain differences.
Candidate inspection is non-exhaustive and does not perform live admission calls
or start watches. Verify provenance with configuration history or
[Kubernetes admission audit records](https://kubernetes.io/docs/reference/access-authn-authz/extensible-admission-controllers/#monitoring-admission-webhooks).

Current-resource warning events in workload `diagnose` and `get_resource(include=events)` exclude a known `involvedObject.uid` mismatch with the already-read resource or its resolved current Pods. Identity still includes the exact API group and scoped namespace. UID-less events remain identity-matched evidence and cannot prove an incarnation. `get_events`, resource history and timeline remain historical views; this filter does not erase earlier incarnations from those surfaces.
61 changes: 46 additions & 15 deletions internal/k8s/workload_pods.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,26 @@ import (
// DaemonSet, ReplicaSet, Job and Workflow directly. An unavailable lister
// returns ErrWorkloadAccessDenied and one still filling its initial sync
// returns ErrWorkloadCacheWarming, so an empty answer is never mistaken for
// "no pods". Pods are sorted by name.
// "no pods". Pods are sorted by name. This name-based API does not identify
// a workload incarnation; callers holding a current root object should use
// WorkloadPodsForUID.
func WorkloadPods(cache *ResourceCache, kind, namespace, name string) ([]*corev1.Pod, error) {
return workloadPods(cache, kind, namespace, name, nil)
}

// WorkloadPodsForUID follows the current root and every intermediate controller
// by API group, kind, name and UID. It reads only already-watched cache listers.
// A missing root UID cannot establish current membership and returns no Pods.
func WorkloadPodsForUID(cache *ResourceCache, kind, namespace, name, group string, uid types.UID) ([]*corev1.Pod, error) {
return workloadPods(cache, kind, namespace, name, &workloadOwnerIdentity{group: group, uid: uid})
}

type workloadOwnerIdentity struct {
group string
uid types.UID
}

func workloadPods(cache *ResourceCache, kind, namespace, name string, identity *workloadOwnerIdentity) ([]*corev1.Pod, error) {
canonical := CanonicalWorkloadKind(kind)
if canonical == "" {
return nil, fmt.Errorf("unsupported workload kind: %s", kind)
Expand All @@ -39,7 +57,7 @@ func WorkloadPods(cache *ResourceCache, kind, namespace, name string) ([]*corev1
if err := requireCovers(cache, "pods", namespace); err != nil {
return nil, err
}
ownedBy, err := workloadOwnershipTest(cache, kind, namespace, name)
ownedBy, err := workloadOwnershipTest(cache, kind, namespace, name, identity)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -100,11 +118,14 @@ func CanonicalWorkloadKind(kind string) string {
// ends at kind/name". Intermediate owners are looked up in the cache; when
// that lister is unavailable the answer is ErrWorkloadAccessDenied rather
// than a silently empty set.
func workloadOwnershipTest(cache *ResourceCache, kind, namespace, name string) (func(*corev1.Pod) bool, error) {
direct := func(pod *corev1.Pod) bool {
owner := metav1.GetControllerOf(pod)
return owner != nil && owner.Kind == kind && owner.Name == name
func workloadOwnershipTest(cache *ResourceCache, kind, namespace, name string, identity *workloadOwnerIdentity) (func(*corev1.Pod) bool, error) {
matchesRoot := func(owner *metav1.OwnerReference) bool {
if identity == nil {
return owner != nil && owner.Kind == kind && owner.Name == name
}
return controllerMatchesIdentity(owner, kind, name, identity.group, identity.uid)
}
direct := func(pod *corev1.Pod) bool { return matchesRoot(metav1.GetControllerOf(pod)) }
switch kind {
case "Deployment", "Rollout":
rsLister := cache.ReplicaSets()
Expand All @@ -123,8 +144,10 @@ func workloadOwnershipTest(cache *ResourceCache, kind, namespace, name string) (
if err != nil {
return false
}
rsOwner := metav1.GetControllerOf(rs)
return rsOwner != nil && rsOwner.Kind == kind && rsOwner.Name == name
if identity != nil && !controllerMatchesIdentity(owner, "ReplicaSet", rs.Name, "apps", rs.UID) {
return false
}
return matchesRoot(metav1.GetControllerOf(rs))
}, nil
case "CronJob":
jobLister := cache.Jobs()
Expand All @@ -143,8 +166,10 @@ func workloadOwnershipTest(cache *ResourceCache, kind, namespace, name string) (
if err != nil {
return false
}
jobOwner := metav1.GetControllerOf(job)
return jobOwner != nil && jobOwner.Kind == "CronJob" && jobOwner.Name == name
if identity != nil && !controllerMatchesIdentity(owner, "Job", job.Name, "batch", job.UID) {
return false
}
return matchesRoot(metav1.GetControllerOf(job))
}, nil
default:
return direct, nil
Expand Down Expand Up @@ -200,13 +225,19 @@ func DirectlyOwnedPods(cache *ResourceCache, namespace, name, group, kind string
continue
}
owner := metav1.GetControllerOf(pod)
if owner == nil || owner.UID != uid || owner.Kind != kind || owner.Name != name {
continue
}
gv, err := schema.ParseGroupVersion(owner.APIVersion)
if err == nil && gv.Group == group {
if controllerMatchesIdentity(owner, kind, name, group, uid) {
result = append(result, pod)
}
}
sort.Slice(result, func(i, j int) bool { return result[i].Name < result[j].Name })
return result, nil
}

// controllerMatchesIdentity is shared by direct and multi-hop ownership checks.
func controllerMatchesIdentity(owner *metav1.OwnerReference, kind, name, group string, uid types.UID) bool {
if owner == nil || uid == "" || owner.UID != uid || owner.Kind != kind || owner.Name != name {
return false
}
gv, err := schema.ParseGroupVersion(owner.APIVersion)
return err == nil && gv.Group == group
}
104 changes: 104 additions & 0 deletions internal/k8s/workload_pods_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/types"
"k8s.io/client-go/kubernetes/fake"

"github.com/skyhook-io/radar/pkg/k8score"
Expand Down Expand Up @@ -276,3 +277,106 @@ func TestDirectlyOwnedPodsRefusesUncoveredNamespace(t *testing.T) {
t.Fatalf("error=%v", err)
}
}

func TestWorkloadPodsForUIDVerifiesEveryControllerHop(t *testing.T) {
for _, tc := range []struct{ kind, group, childKind, childAPI string }{
{"Deployment", "apps", "ReplicaSet", "apps/v1"},
{"Rollout", "argoproj.io", "ReplicaSet", "apps/v1"},
{"CronJob", "batch", "Job", "batch/v1"},
} {
t.Run(tc.kind, func(t *testing.T) {
controller := true
ref := func(api, kind, name, uid string) metav1.OwnerReference {
return metav1.OwnerReference{APIVersion: api, Kind: kind, Name: name, UID: types.UID(uid), Controller: &controller}
}
root := ref(tc.group+"/v1", tc.kind, "fleet", "current-root")
oldRoot := root
oldRoot.UID = "previous-root"
foreignRoot := root
foreignRoot.APIVersion = "other.example/v1"
objects := []runtime.Object{}
watches := map[string]bool{k8score.Pods: true}
for _, child := range []struct {
name, uid string
owner metav1.OwnerReference
}{
{"active", "active-child", root}, {"previous", "previous-child", oldRoot}, {"foreign", "foreign-child", foreignRoot},
} {
metadata := metav1.ObjectMeta{Namespace: "shop", Name: child.name, UID: types.UID(child.uid), OwnerReferences: []metav1.OwnerReference{child.owner}}
if tc.childKind == "ReplicaSet" {
objects = append(objects, &appsv1.ReplicaSet{ObjectMeta: metadata})
watches[k8score.ReplicaSets] = true
} else {
objects = append(objects, &batchv1.Job{ObjectMeta: metadata})
watches[k8score.Jobs] = true
}
owner := ref(tc.childAPI, tc.childKind, child.name, child.uid)
objects = append(objects, ownedPod("shop", child.name+"-pod", &owner, nil))
}
staleChild := ref(tc.childAPI, tc.childKind, "active", "replaced-child")
foreignChild := ref("other.example/v1", tc.childKind, "active", "active-child")
objects = append(objects, ownedPod("shop", "stale-child-pod", &staleChild, nil), ownedPod("shop", "foreign-child-pod", &foreignChild, nil))
cache := newOwnershipCache(t, watches, objects...)
got, err := WorkloadPodsForUID(cache, tc.kind, "shop", "fleet", tc.group, "current-root")
if err != nil {
t.Fatal(err)
}
if len(got) != 1 || got[0].Name != "active-pod" {
t.Fatalf("current chain pods = %v", got)
}
got, err = WorkloadPodsForUID(cache, tc.kind, "shop", "fleet", tc.group, "")
if err != nil || len(got) != 0 {
t.Fatalf("missing root UID: pods=%v err=%v", got, err)
}
})
}
}

func TestWorkloadPodsForUIDVerifiesDirectControllers(t *testing.T) {
for _, tc := range []struct{ kind, api, group string }{
{"StatefulSet", "apps/v1", "apps"}, {"DaemonSet", "apps/v1", "apps"}, {"ReplicaSet", "apps/v1", "apps"}, {"Job", "batch/v1", "batch"}, {"Workflow", "argoproj.io/v1alpha1", "argoproj.io"},
} {
t.Run(tc.kind, func(t *testing.T) {
controller := true
current := metav1.OwnerReference{APIVersion: tc.api, Kind: tc.kind, Name: "fleet", UID: "current", Controller: &controller}
old := current
old.UID = "previous"
foreign := current
foreign.APIVersion = "other.example/v1"
cache := newOwnershipCache(t, map[string]bool{k8score.Pods: true}, ownedPod("shop", "current-pod", &current, nil), ownedPod("shop", "old-pod", &old, nil), ownedPod("shop", "foreign-pod", &foreign, nil))
got, err := WorkloadPodsForUID(cache, tc.kind, "shop", "fleet", tc.group, "current")
if err != nil || len(got) != 1 || got[0].Name != "current-pod" {
t.Fatalf("pods=%v err=%v", got, err)
}
})
}
}

func TestWorkloadPodsForUIDUsesOnlyReadyCache(t *testing.T) {
controller := true
rs := &appsv1.ReplicaSet{ObjectMeta: metav1.ObjectMeta{Name: "rs", Namespace: "shop", UID: "rs-uid", OwnerReferences: []metav1.OwnerReference{{APIVersion: "apps/v1", Kind: "Deployment", Name: "fleet", UID: "root-uid", Controller: &controller}}}}
owner := metav1.OwnerReference{APIVersion: "apps/v1", Kind: "ReplicaSet", Name: "rs", UID: rs.UID, Controller: &controller}
client := fake.NewClientset(rs, ownedPod("shop", "current", &owner, nil))
core, err := k8score.NewResourceCache(k8score.CacheConfig{Client: client, ResourceTypes: map[string]bool{k8score.Pods: true, k8score.ReplicaSets: true}, DeferredTypes: map[string]bool{}})
if err != nil {
t.Fatal(err)
}
defer core.Stop()
cache := &ResourceCache{ResourceCache: core}
before := len(client.Actions())
for range 100 {
got, err := WorkloadPodsForUID(cache, "deployment", "shop", "fleet", "apps", "root-uid")
if err != nil || len(got) != 1 {
t.Fatalf("pods=%v err=%v", got, err)
}
}
for _, action := range client.Actions()[before:] {
if action.GetVerb() == "get" || action.GetVerb() == "list" {
t.Fatalf("hot relationship lookup fetched %s %s", action.GetVerb(), action.GetResource())
}
}
denied := newOwnershipCache(t, map[string]bool{k8score.Pods: true})
if _, err := WorkloadPodsForUID(denied, "deployment", "shop", "fleet", "apps", "root-uid"); !errors.Is(err, ErrWorkloadAccessDenied) {
t.Fatalf("unavailable hop became empty membership: %v", err)
}
}
Loading
Loading