diff --git a/docs/mcp.md b/docs/mcp.md index fbe55dde71..1b83d8d749 100644 --- a/docs/mcp.md +++ b/docs/mcp.md @@ -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. diff --git a/internal/k8s/workload_pods.go b/internal/k8s/workload_pods.go index 7dff99dfce..64291b28d2 100644 --- a/internal/k8s/workload_pods.go +++ b/internal/k8s/workload_pods.go @@ -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) @@ -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 } @@ -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() @@ -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() @@ -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 @@ -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 +} diff --git a/internal/k8s/workload_pods_test.go b/internal/k8s/workload_pods_test.go index 2f1b1a55f5..1ec89a62c1 100644 --- a/internal/k8s/workload_pods_test.go +++ b/internal/k8s/workload_pods_test.go @@ -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" @@ -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", ¤t, 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) + } +} diff --git a/internal/mcp/events_tool_test.go b/internal/mcp/events_tool_test.go index 12ec47e0ee..2b82aa8f8f 100644 --- a/internal/mcp/events_tool_test.go +++ b/internal/mcp/events_tool_test.go @@ -9,9 +9,11 @@ import ( "time" "github.com/modelcontextprotocol/go-sdk/mcp" + appsv1 "k8s.io/api/apps/v1" 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/internal/k8s" @@ -178,7 +180,7 @@ func TestAttachResourceExtras_EventsTotalGroups(t *testing.T) { var result map[string]any for time.Now().Before(deadline) { result = map[string]any{} - attachResourceExtras(context.Background(), cache, result, map[string]bool{"events": true}, "deployment", "apps", "shop", "web") + attachResourceExtras(context.Background(), cache, result, map[string]bool{"events": true}, "deployment", "apps", "shop", "web", "") if evs, ok := result["events"].([]aicontext.DeduplicatedEvent); ok && len(evs) == 10 { break } @@ -195,7 +197,7 @@ func TestAttachResourceExtras_EventsTotalGroups(t *testing.T) { // Under the cap: no truncation field. few := map[string]any{} - attachResourceExtras(context.Background(), cache, few, map[string]bool{"events": true}, "deployment", "apps", "shop", "missing") + attachResourceExtras(context.Background(), cache, few, map[string]bool{"events": true}, "deployment", "apps", "shop", "missing", "") if _, present := few["eventsTotalGroups"]; present { t.Errorf("eventsTotalGroups present with no truncation: %+v", few) } @@ -227,7 +229,7 @@ func TestFetchEventsForResource_ReportsPreCapGroupTotal(t *testing.T) { err error ) for time.Now().Before(deadline) { - groups, total, err = fetchEventsForResource(k8s.GetResourceCache(), "deployments", "apps", "shop", "web", nil, 10) + groups, total, err = fetchEventsForResource(k8s.GetResourceCache(), "deployments", "apps", "shop", "web", "", nil, 10) if err != nil || total == 12 { break } @@ -240,7 +242,7 @@ func TestFetchEventsForResource_ReportsPreCapGroupTotal(t *testing.T) { t.Fatalf("groups=%d total=%d, want capped response of 10 from 12 deduplicated groups", len(groups), total) } - groups, total, err = fetchEventsForResource(k8s.GetResourceCache(), "deployments", "apps", "shop", "missing", nil, 10) + groups, total, err = fetchEventsForResource(k8s.GetResourceCache(), "deployments", "apps", "shop", "missing", "", nil, 10) if err != nil || len(groups) != 0 || total != 0 { t.Fatalf("missing resource groups=%d total=%d err=%v, want empty successful result", len(groups), total, err) } @@ -278,7 +280,7 @@ func TestFetchEventsForResource_RolloutWarningSurvivesKindAndGroupFiltering(t *t err error ) for time.Now().Before(deadline) { - groups, total, err = fetchEventsForResource(k8s.GetResourceCache(), "rollouts", "argoproj.io", "shop", "checkout", nil, 10) + groups, total, err = fetchEventsForResource(k8s.GetResourceCache(), "rollouts", "argoproj.io", "shop", "checkout", "", nil, 10) if err != nil || total == 1 { break } @@ -298,12 +300,156 @@ func TestFilterEventsByInvolvedObject_RequiresExactAPIGroup(t *testing.T) { {Type: corev1.EventTypeWarning, Reason: "Knative", InvolvedObject: corev1.ObjectReference{APIVersion: "serving.knative.dev/v1", Kind: "Service", Name: "api"}}, } - core := filterEventsByInvolvedObject(events, "Service", "", "api", nil) + core := filterEventsByInvolvedObject(events, "Service", "", "api", "", nil) if len(core) != 1 || core[0].Reason != "Core" { t.Fatalf("core Service events = %+v, want only core/v1", core) } - knative := filterEventsByInvolvedObject(events, "Service", "serving.knative.dev", "api", nil) + knative := filterEventsByInvolvedObject(events, "Service", "serving.knative.dev", "api", "", nil) if len(knative) != 1 || knative[0].Reason != "Knative" { t.Fatalf("Knative Service events = %+v, want only serving.knative.dev", knative) } } + +func TestFilterEventsByInvolvedObject_CurrentIncarnations(t *testing.T) { + event := func(reason, kind, group, name, uid string) *corev1.Event { + version := "v1" + if group != "" { + version = group + "/v1" + } + return &corev1.Event{Type: corev1.EventTypeWarning, Reason: reason, InvolvedObject: corev1.ObjectReference{APIVersion: version, Kind: kind, Name: name, UID: types.UID(uid)}} + } + events := []*corev1.Event{ + event("CurrentRoot", "Deployment", "apps", "web", "root-now"), + event("ReplacedRoot", "Deployment", "apps", "web", "root-before"), + event("UIDlessRoot", "Deployment", "apps", "web", ""), + event("CurrentPod", "Pod", "", "worker", "pod-now"), + event("ReplacedPod", "Pod", "", "worker", "pod-before"), + event("UIDlessPod", "Pod", "", "worker", ""), + event("WrongGroup", "Deployment", "custom.example.io", "web", "root-now"), + event("WrongPodGroup", "Pod", "custom.example.io", "worker", "pod-now"), + event("UnknownPod", "Pod", "", "foreign", "other"), + } + got := filterEventsByInvolvedObject(events, "Deployment", "apps", "web", "root-now", map[string]types.UID{"worker": "pod-now"}) + var reasons []string + for _, e := range got { + reasons = append(reasons, e.Reason) + } + if strings.Join(reasons, ",") != "CurrentRoot,UIDlessRoot,CurrentPod,UIDlessPod" { + t.Fatalf("current evidence = %v", reasons) + } + // Supplemental includes do not enroll workload Pod events. + if got := filterEventsByInvolvedObject(events, "Deployment", "apps", "web", "root-now", nil); len(got) != 2 { + t.Fatalf("root-only evidence = %v", got) + } + // No Node UID=name exception: actual Kubernetes identity wins. + nodes := []*corev1.Event{event("CurrentNode", "Node", "", "node-a", "node-uid"), event("NameUID", "Node", "", "node-a", "node-a")} + if got := filterEventsByInvolvedObject(nodes, "Node", "", "node-a", "node-uid", nil); len(got) != 1 || got[0].Reason != "CurrentNode" { + t.Fatalf("Node evidence = %v", got) + } +} + +func TestGetResourceEvents_UsesObservedResourceUID(t *testing.T) { + defer k8s.ResetTestState() + now := metav1.Now() + root := &appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: "web", Namespace: "shop", UID: "root-now"}} + events := []runtime.Object{root} + for _, uid := range []types.UID{"root-now", "root-before"} { + events = append(events, &corev1.Event{ObjectMeta: metav1.ObjectMeta{Name: string(uid), Namespace: "shop"}, Type: corev1.EventTypeWarning, Reason: string(uid), Message: string(uid), LastTimestamp: now, InvolvedObject: corev1.ObjectReference{APIVersion: "apps/v1", Kind: "Deployment", Namespace: "shop", Name: "web", UID: uid}}) + } + client := fake.NewSimpleClientset(events...) + if err := k8s.InitTestResourceCache(client); err != nil { + t.Fatal(err) + } + var response struct { + Events []aicontext.DeduplicatedEvent `json:"events"` + } + deadline := time.Now().Add(2 * time.Second) + for { + result, _, err := handleGetResource(t.Context(), nil, getResourceInput{Kind: "deployment", Namespace: "shop", Name: "web", Include: "events", Context: "none"}) + if err == nil { + if err := json.Unmarshal([]byte(extractText(t, result)), &response); err != nil { + t.Fatal(err) + } + if len(response.Events) > 0 { + break + } + } + if time.Now().After(deadline) { + t.Fatalf("events never loaded: %+v err=%v", response, err) + } + time.Sleep(10 * time.Millisecond) + } + if len(response.Events) != 1 || response.Events[0].Reason != "root-now" { + t.Fatalf("get_resource current evidence = %+v", response.Events) + } + pods := []*corev1.Pod{{ObjectMeta: metav1.ObjectMeta{Name: "worker", Namespace: "shop", UID: "pod-now"}}} + for _, uid := range []types.UID{"pod-now", "pod-before"} { + _, err := client.CoreV1().Events("shop").Create(t.Context(), &corev1.Event{ObjectMeta: metav1.ObjectMeta{Name: string(uid), Namespace: "shop"}, Type: corev1.EventTypeWarning, Reason: string(uid), Message: string(uid), LastTimestamp: now, InvolvedObject: corev1.ObjectReference{APIVersion: "v1", Kind: "Pod", Namespace: "shop", Name: "worker", UID: uid}}, metav1.CreateOptions{}) + if err != nil { + t.Fatal(err) + } + } + deadline = time.Now().Add(2 * time.Second) + for { + groups, total, err := fetchEventsForResource(k8s.GetResourceCache(), "deployment", "apps", "shop", "web", root.UID, pods, 10) + if err != nil { + t.Fatal(err) + } + if total == 2 { + for _, g := range groups { + if g.Reason != "root-now" && g.Reason != "pod-now" { + t.Fatalf("diagnose current evidence = %+v", groups) + } + } + break + } + if time.Now().After(deadline) { + t.Fatalf("diagnose event groups = %+v", groups) + } + time.Sleep(10 * time.Millisecond) + } +} + +func TestDiagnoseExcludesPodsFromPreviousWorkloadAndIntermediateOwners(t *testing.T) { + defer k8s.ResetTestState() + controller := true + root := &appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: "fleet", Namespace: "shop", UID: "current-root"}, Spec: appsv1.DeploymentSpec{Selector: &metav1.LabelSelector{MatchLabels: map[string]string{"app": "fleet"}}, Template: corev1.PodTemplateSpec{ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{"app": "fleet"}}, Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "app"}}}}}} + objects := []runtime.Object{root} + for _, item := range []struct{ name, childUID, rootUID, podOwnerUID string }{ + {"current", "current-rs", "current-root", "current-rs"}, + {"previous-root", "old-rs", "old-root", "old-rs"}, + {"previous-rs", "replaced-rs", "current-root", "previous-rs"}, + } { + objects = append(objects, &appsv1.ReplicaSet{ObjectMeta: metav1.ObjectMeta{Name: item.name, Namespace: "shop", UID: types.UID(item.childUID), OwnerReferences: []metav1.OwnerReference{{APIVersion: "apps/v1", Kind: "Deployment", Name: "fleet", UID: types.UID(item.rootUID), Controller: &controller}}}}) + pod := &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: item.name, Namespace: "shop", UID: types.UID(item.name + "-uid"), OwnerReferences: []metav1.OwnerReference{{APIVersion: "apps/v1", Kind: "ReplicaSet", Name: item.name, UID: types.UID(item.podOwnerUID), Controller: &controller}}}, Spec: corev1.PodSpec{Containers: []corev1.Container{{Name: "app"}}}} + event := &corev1.Event{ObjectMeta: metav1.ObjectMeta{Name: item.name, Namespace: "shop"}, Type: corev1.EventTypeWarning, Reason: item.name, Message: item.name, LastTimestamp: metav1.Now(), InvolvedObject: corev1.ObjectReference{APIVersion: "v1", Kind: "Pod", Name: pod.Name, Namespace: pod.Namespace, UID: pod.UID}} + objects = append(objects, pod, event) + } + if err := k8s.InitTestResourceCache(fake.NewClientset(objects...)); err != nil { + t.Fatal(err) + } + var response diagnoseResponse + deadline := time.Now().Add(2 * time.Second) + for { + result, _, err := handleDiagnose(t.Context(), nil, testDiagnoseInput("deployment", "shop", "fleet")) + if err != nil { + t.Fatal(err) + } + if err := json.Unmarshal([]byte(extractText(t, result)), &response); err != nil { + t.Fatal(err) + } + if response.Pods > 0 && len(response.Events) > 0 { + break + } + if time.Now().After(deadline) { + t.Fatalf("fixture did not sync: %+v", response) + } + time.Sleep(10 * time.Millisecond) + } + if response.Pods != 1 || len(response.PodNames) != 1 || response.PodNames[0] != "current" { + t.Fatalf("diagnose current pods=%d names=%v", response.Pods, response.PodNames) + } + if len(response.Events) != 1 || response.Events[0].Reason != "current" { + t.Fatalf("diagnose leaked previous controller warnings: %+v", response.Events) + } +} diff --git a/internal/mcp/tools.go b/internal/mcp/tools.go index 8a397f4391..f25cd2d92d 100644 --- a/internal/mcp/tools.go +++ b/internal/mcp/tools.go @@ -16,11 +16,13 @@ import ( "github.com/modelcontextprotocol/go-sdk/mcp" corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" "github.com/skyhook-io/radar/internal/filter" "github.com/skyhook-io/radar/internal/helm" @@ -1198,7 +1200,11 @@ func handleGetResource(ctx context.Context, req *mcp.CallToolRequest, input getR } canonicalGroup = gvk.Group } - attachResourceExtras(ctx, cache, result, includes, canonicalKind, canonicalGroup, namespace, name) + objectMeta, err := meta.Accessor(rawObj) + if err != nil { + return nil, nil, fmt.Errorf("failed to read resource metadata: %w", err) + } + attachResourceExtras(ctx, cache, result, includes, canonicalKind, canonicalGroup, namespace, name, objectMeta.GetUID()) } return toJSONResult(result) } @@ -1296,7 +1302,7 @@ func buildMCPResourceContextWithStaleChecks(ctx context.Context, obj runtime.Obj // attachResourceExtras populates optional extras (events, metrics, logs) on // the result map based on the includes set. relationship synthesis moved to // resourceContext via Build and is no longer routed through this function. -func attachResourceExtras(ctx context.Context, cache *k8s.ResourceCache, result map[string]any, includes map[string]bool, kind, group, namespace, name string) { +func attachResourceExtras(ctx context.Context, cache *k8s.ResourceCache, result map[string]any, includes map[string]bool, kind, group, namespace, name string, uid types.UID) { if includes["events"] { if eventLister := cache.Events(); eventLister != nil { var events []*corev1.Event @@ -1313,9 +1319,9 @@ func attachResourceExtras(ctx context.Context, cache *k8s.ResourceCache, result // Supplemental include — controller-level events only. Pod-level // events on a workload's pods (CrashLoopBackOff, etc.) require // resolving the pod set; that's the diagnose tool's job, not - // this include's. nil podNames intentionally restricts to + // this include's. nil podUIDs intentionally restricts to // InvolvedObject == this kind+name. - matched := filterEventsByInvolvedObject(events, normalizeDisplayKind(kind), group, name, nil) + matched := filterEventsByInvolvedObject(events, normalizeDisplayKind(kind), group, name, uid, nil) if len(matched) > 0 { deduplicated, totalGroups := aicontext.DeduplicateEventsN(matched, 10) result["events"] = deduplicated diff --git a/internal/mcp/tools_diagnose.go b/internal/mcp/tools_diagnose.go index 30d74c2c59..cde69808b8 100644 --- a/internal/mcp/tools_diagnose.go +++ b/internal/mcp/tools_diagnose.go @@ -13,9 +13,11 @@ import ( "github.com/modelcontextprotocol/go-sdk/mcp" appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/meta" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" "github.com/skyhook-io/radar/internal/issues" "github.com/skyhook-io/radar/internal/k8s" @@ -313,6 +315,10 @@ func handleDiagnose(ctx context.Context, _ *mcp.CallToolRequest, input diagnoseI } } k8s.SetTypeMeta(obj) + objectMeta, err := meta.Accessor(obj) + if err != nil { + return nil, nil, fmt.Errorf("failed to read resource metadata: %w", err) + } gvk := obj.GetObjectKind().GroupVersionKind() canonicalGroup := gvk.Group canonicalKind := gvk.Kind @@ -433,7 +439,7 @@ func handleDiagnose(ctx context.Context, _ *mcp.CallToolRequest, input diagnoseI ) } - events, eventsTotalGroups, eventsErr := fetchEventsForResource(cache, kindNorm, canonicalGroup, input.Namespace, input.Name, pods, 10) + events, eventsTotalGroups, eventsErr := fetchEventsForResource(cache, kindNorm, canonicalGroup, input.Namespace, input.Name, objectMeta.GetUID(), pods, 10) resp.Events = events resp.EventsTotalGroups = eventsTotalGroups if eventsErr != nil { @@ -821,8 +827,8 @@ func workloadDiagnoseGroup(kind string) string { // resolveDiagnosePods returns the set of pods to fetch logs from. For // kind=pods that's just the requested pod; for workload kinds it resolves -// via the workload's pod selector and the cache's pod-by-workload index. -func resolveDiagnosePods(cache *k8s.ResourceCache, kindNorm, namespace, name string, obj any) ([]*corev1.Pod, error) { +// via the current root UID and the cached controller ownership chain. +func resolveDiagnosePods(cache *k8s.ResourceCache, kindNorm, namespace, name string, obj runtime.Object) ([]*corev1.Pod, error) { if kindNorm == "pods" { pod, ok := obj.(*corev1.Pod) if !ok || pod == nil { @@ -830,11 +836,14 @@ func resolveDiagnosePods(cache *k8s.ResourceCache, kindNorm, namespace, name str } return []*corev1.Pod{pod}, nil } - // Membership is controller ownership, the relation kube-state-metrics - // records too, so the bundle, its vitals and the workload page name the - // same pods; a selector would also match bare pods and a sibling - // controller's pods during a Rollout migration. - pods, err := k8s.WorkloadPods(cache, kindNorm, namespace, name) + // Qualify the full controller chain by the root already read for this + // diagnosis. Same-name previous roots and replaced intermediate owners + // must not contribute logs, vitals or warning events. + objectMeta, err := meta.Accessor(obj) + if err != nil { + return nil, fmt.Errorf("failed to read workload metadata: %w", err) + } + pods, err := k8s.WorkloadPodsForUID(cache, kindNorm, namespace, name, obj.GetObjectKind().GroupVersionKind().Group, objectMeta.GetUID()) if errors.Is(err, k8s.ErrWorkloadCacheWarming) { return nil, errPodsCacheWarming } @@ -865,7 +874,7 @@ var errPodsCacheWarming = errors.New("the pod cache is still loading") // "no warnings exist" from "apiserver list failed and we couldn't tell" // — diagnose surfaces it as EventsError so the agent doesn't read empty // events as ground truth. -func fetchEventsForResource(cache *k8s.ResourceCache, kind, group, namespace, name string, pods []*corev1.Pod, limit int) ([]aicontext.DeduplicatedEvent, int, error) { +func fetchEventsForResource(cache *k8s.ResourceCache, kind, group, namespace, name string, uid types.UID, pods []*corev1.Pod, limit int) ([]aicontext.DeduplicatedEvent, int, error) { eventLister := cache.Events() if eventLister == nil { // Mirror attachResourceExtras / get_resource(include=events): surface @@ -878,13 +887,13 @@ func fetchEventsForResource(cache *k8s.ResourceCache, kind, group, namespace, na log.Printf("[mcp] diagnose: failed to list events for %s/%s/%s: %v", kind, namespace, name, err) return nil, 0, err } - podNames := make(map[string]bool, len(pods)) + podUIDs := make(map[string]types.UID, len(pods)) for _, p := range pods { if p != nil { - podNames[p.Name] = true + podUIDs[p.Name] = p.UID } } - matched := filterEventsByInvolvedObject(events, normalizeDisplayKind(kind), group, name, podNames) + matched := filterEventsByInvolvedObject(events, normalizeDisplayKind(kind), group, name, uid, podUIDs) if len(matched) == 0 { return nil, 0, nil } @@ -894,18 +903,6 @@ func fetchEventsForResource(cache *k8s.ResourceCache, kind, group, namespace, na return dedup, totalGroups, nil } -// filterEventsByInvolvedObject keeps Warning events whose InvolvedObject -// matches either the controller (displayKind+name) OR any of the pods in -// podNames (skipped when displayKind is "Pod" — the controller branch -// above already covers single-pod and otherwise this branch would -// double-count). -// -// Filters to Type==Warning intentionally — the diagnose tool description -// + get_resource(include=events) both promise warning events only. -// Normal events (Pulled / Created / Scheduled) would pollute triage by -// reading as "things worth diagnosing" when they're just lifecycle -// breadcrumbs. -// // networkTraceKind returns the canonical entry-kind name for trace // diagnostics, accepting plural and lowercase forms the way agents tend to // write them. Empty string + false means "not a network entry kind" - the @@ -1130,11 +1127,12 @@ func handleNetworkTraceDiagnose(ctx context.Context, input diagnoseInput, kind s return toJSONResult(resp) } -// Shared between diagnose (passes resolved pod names for full workload -// coverage) and attachResourceExtras / get_resource include=events -// (passes nil — supplemental fetch; callers wanting pod-level events should -// use the diagnose tool which does the workload→pods resolution). -func filterEventsByInvolvedObject(events []*corev1.Event, displayKind, group, name string, podNames map[string]bool) []corev1.Event { +// filterEventsByInvolvedObject keeps Warning events for the current resource +// or its resolved Pods. A known UID mismatch rejects a previous incarnation; +// UID-less events remain identity-matched evidence without proof of incarnation. +// Namespace is scoped by the event lister. get_resource passes no Pods; diagnose +// supplies its authorized current Pod set. Normal lifecycle events stay out. +func filterEventsByInvolvedObject(events []*corev1.Event, displayKind, group, name string, uid types.UID, podUIDs map[string]types.UID) []corev1.Event { var matched []corev1.Event for _, e := range events { if e.Type != corev1.EventTypeWarning { @@ -1144,12 +1142,15 @@ func filterEventsByInvolvedObject(events []*corev1.Event, displayKind, group, na if e.InvolvedObject.APIVersion == "" { involvedGroup = resourceid.GroupForBuiltinKind(e.InvolvedObject.Kind) } - if strings.EqualFold(e.InvolvedObject.Kind, displayKind) && strings.EqualFold(involvedGroup, group) && e.InvolvedObject.Name == name { + if strings.EqualFold(e.InvolvedObject.Kind, displayKind) && strings.EqualFold(involvedGroup, group) && e.InvolvedObject.Name == name && (uid == "" || e.InvolvedObject.UID == "" || e.InvolvedObject.UID == uid) { matched = append(matched, *e) continue } - if displayKind != "Pod" && involvedGroup == "" && strings.EqualFold(e.InvolvedObject.Kind, "Pod") && podNames[e.InvolvedObject.Name] { - matched = append(matched, *e) + if displayKind != "Pod" && involvedGroup == "" && strings.EqualFold(e.InvolvedObject.Kind, "Pod") { + podUID, ok := podUIDs[e.InvolvedObject.Name] + if ok && (podUID == "" || e.InvolvedObject.UID == "" || e.InvolvedObject.UID == podUID) { + matched = append(matched, *e) + } } } return matched diff --git a/internal/mcp/tools_diagnose_test.go b/internal/mcp/tools_diagnose_test.go index 4d69793f8d..e1fdd1207a 100644 --- a/internal/mcp/tools_diagnose_test.go +++ b/internal/mcp/tools_diagnose_test.go @@ -34,17 +34,17 @@ func setupFakeCacheForDiagnoseTests(t *testing.T) { // Pods belong to a workload by controller ownership, so the fixture // carries the Deployment → ReplicaSet → Pod chain a real cluster has. isController := true - rsOwner := metav1.OwnerReference{APIVersion: "apps/v1", Kind: "Deployment", Name: deployName, Controller: &isController} - podOwner := metav1.OwnerReference{APIVersion: "apps/v1", Kind: "ReplicaSet", Name: "cart-abc", Controller: &isController} + rsOwner := metav1.OwnerReference{APIVersion: "apps/v1", Kind: "Deployment", Name: deployName, UID: "current-cart", Controller: &isController} + podOwner := metav1.OwnerReference{APIVersion: "apps/v1", Kind: "ReplicaSet", Name: "cart-abc", UID: "current-cart-rs", Controller: &isController} fakeClient := fake.NewClientset( &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: ns}, Status: corev1.NamespaceStatus{Phase: corev1.NamespaceActive}}, &appsv1.ReplicaSet{ - ObjectMeta: metav1.ObjectMeta{Name: "cart-abc", Namespace: ns, Labels: selector, OwnerReferences: []metav1.OwnerReference{rsOwner}}, + ObjectMeta: metav1.ObjectMeta{Name: "cart-abc", Namespace: ns, UID: "current-cart-rs", Labels: selector, OwnerReferences: []metav1.OwnerReference{rsOwner}}, Spec: appsv1.ReplicaSetSpec{Selector: &metav1.LabelSelector{MatchLabels: selector}}, }, &appsv1.Deployment{ - ObjectMeta: metav1.ObjectMeta{Name: deployName, Namespace: ns}, + ObjectMeta: metav1.ObjectMeta{Name: deployName, Namespace: ns, UID: "current-cart"}, Spec: appsv1.DeploymentSpec{ Selector: &metav1.LabelSelector{MatchLabels: selector}, Template: corev1.PodTemplateSpec{ @@ -355,10 +355,10 @@ func TestHandleDiagnose_PodNotFound(t *testing.T) { } // TestHandleDiagnose_DeploymentResolvesPods exercises the workload-rooted -// path (kind=deployment → workload selector → fan-out to matching pods), +// path (kind=deployment → cached UID-qualified controller chain → fan-out to owned pods), // which is the diagnose tool's headline use case. The pod-only tests above // never traverse this branch — without this test, a regression in -// GetWorkloadSelector / GetPodsForWorkload / selector matching would ship +// controller-chain matching would ship // undetected on the most common debug journey ("CrashLoopBackOff on a // Deployment"). The fake test environment has no kube client on ctx, so // logs surface as LogsError rather than empty arrays — that's the @@ -381,9 +381,9 @@ func TestHandleDiagnose_DeploymentResolvesPods(t *testing.T) { if !strings.Contains(body, `"name":"cart"`) { t.Errorf("expected deployment name in response: %s", body) } - // Selector resolution should find the matching pod. + // Current controller ownership should find exactly one Pod. if !strings.Contains(body, `"pods":1`) { - t.Errorf("expected pods:1 (selector matched 1 pod): %s", body) + t.Errorf("expected pods:1 (current controller owns 1 pod): %s", body) } // No kube client on ctx in tests — diagnose surfaces this distinctly. if !strings.Contains(body, "logsError") { diff --git a/internal/mcp/tools_rollouts_test.go b/internal/mcp/tools_rollouts_test.go index 111f0fee7a..437c5f899f 100644 --- a/internal/mcp/tools_rollouts_test.go +++ b/internal/mcp/tools_rollouts_test.go @@ -58,7 +58,7 @@ func TestRevisionCapableKind(t *testing.T) { // reports it as an unknown include alongside the data it just attached. func TestRevisionsIsAKnownIncludeToken(t *testing.T) { result := map[string]any{} - attachResourceExtras(t.Context(), nil, result, map[string]bool{"revisions": true}, "pod", "", "default", "web") + attachResourceExtras(t.Context(), nil, result, map[string]bool{"revisions": true}, "pod", "", "default", "web", "") if msg, present := result["includeError"]; present { t.Errorf("revisions reported as unknown include: %v", msg)