Skip to content

Commit 8351a67

Browse files
Colvin-Ycodex
andcommitted
fix(diagnostic): Scope Pod CSI topology
Restrict Pod storage diagnostics to the relevant CSI node agent while preserving controller evidence. Allow diagnostic graph loading for cluster-scoped storage resources without sending stale namespaces. Co-Authored-By: Codex <codex@openai.com>
1 parent 84f47d9 commit 8351a67

14 files changed

Lines changed: 416 additions & 24 deletions

File tree

AI_CONTRACT.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -251,7 +251,7 @@ It should not be used as the only signal for graph facts already present in node
251251
- this pod mounts PVC X, which is bound to PV Y
252252
- this PV uses CSI driver X
253253
- this PV is managed by configured CSI controller components
254-
- no configured CSI node agent was found for the PV affinity node in the current observed graph slice
254+
- no configured CSI node agent was found for the PV affinity or consuming Pod node in the current observed graph slice
255255
- this workload owns the pod through the owner chain
256256

257257
### Unsafe conclusions

cmd/kubernetes-ontology-viewer/main.go

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -157,7 +157,10 @@ func (h *handler) serveDiagnostic(w http.ResponseWriter, r *http.Request) {
157157
kind := first(q, "kind", "Pod")
158158
namespace := first(q, "namespace", "")
159159
name := first(q, "name", "")
160-
if namespace == "" && (strings.EqualFold(kind, "Pod") || strings.EqualFold(kind, "Workload")) {
160+
if isDiagnosticClusterScopedKind(kind) {
161+
namespace = ""
162+
}
163+
if namespace == "" && isDiagnosticNamespacedKind(kind) {
161164
writeJSON(w, map[string]string{"error": "namespace is required"}, http.StatusBadRequest)
162165
return
163166
}
@@ -189,6 +192,24 @@ func (h *handler) serveDiagnostic(w http.ResponseWriter, r *http.Request) {
189192
writeJSON(w, data, http.StatusOK)
190193
}
191194

195+
func isDiagnosticNamespacedKind(kind string) bool {
196+
switch strings.ToLower(strings.TrimSpace(kind)) {
197+
case "pod", "workload", "pvc", "persistentvolumeclaim":
198+
return true
199+
default:
200+
return false
201+
}
202+
}
203+
204+
func isDiagnosticClusterScopedKind(kind string) bool {
205+
switch strings.ToLower(strings.TrimSpace(kind)) {
206+
case "pv", "persistentvolume", "storageclass", "csidriver":
207+
return true
208+
default:
209+
return false
210+
}
211+
}
212+
192213
func (h *handler) serveExpand(w http.ResponseWriter, r *http.Request) {
193214
q := r.URL.Query()
194215
server := first(q, "server", h.defaultServer)

docs/ontology/kubernetes-ontology.owl

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -422,7 +422,7 @@
422422
</owl:ObjectProperty>
423423
<owl:ObjectProperty rdf:about="#served_by_csi_node_agent">
424424
<rdfs:label>served_by_csi_node_agent</rdfs:label>
425-
<rdfs:comment>Relates a PersistentVolume to the CSI node-agent Pod on its affinity node.</rdfs:comment>
425+
<rdfs:comment>Relates a PersistentVolume to a CSI node-agent Pod on a PV affinity node or consuming Pod node.</rdfs:comment>
426426
<ko:edgeKind>served_by_csi_node_agent</ko:edgeKind>
427427
<rdfs:domain rdf:resource="#PV"/>
428428
<rdfs:range rdf:resource="#Pod"/>

internal/graph/builder.go

Lines changed: 40 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package graph
22

33
import (
4+
"sort"
45
"strings"
56

67
"github.com/Colvin-Y/kubernetes-ontology/internal/collect/k8s"
@@ -313,12 +314,11 @@ func (b *Builder) Build(snapshot k8s.Snapshot) ([]model.Node, []model.Edge) {
313314
if !ok {
314315
continue
315316
}
316-
affinityNodeName, _ := pv.CSI["nodeAffinity"]
317317
pvNode, ok := findNode(nodes, pvID)
318318
if !ok {
319319
continue
320320
}
321-
correlation := correlator.Correlate(pvNode, affinityNodeName, infraPods)
321+
correlation := infer.CorrelatePVToCSIComponents(correlator, pvNode, csiAgentNodeNamesForPV(pv, snapshot), infraPods)
322322
edges = append(edges, correlation.Edges...)
323323
b.lastEvidence = append(b.lastEvidence, correlation.Evidence...)
324324
}
@@ -374,6 +374,44 @@ func storageClassNameForPVC(pvc k8sresources.PVC, pvs []k8sresources.PV) string
374374
return ""
375375
}
376376

377+
func csiAgentNodeNamesForPV(pv k8sresources.PV, snapshot k8s.Snapshot) []string {
378+
seen := make(map[string]struct{})
379+
addNode := func(nodeName string) {
380+
nodeName = strings.TrimSpace(nodeName)
381+
if nodeName == "" {
382+
return
383+
}
384+
seen[nodeName] = struct{}{}
385+
}
386+
addNode(pv.CSI["nodeAffinity"])
387+
for _, pvc := range snapshot.PVCs {
388+
if pvc.VolumeName != pv.Metadata.Name {
389+
continue
390+
}
391+
for _, pod := range snapshot.Pods {
392+
if pod.Metadata.Namespace != pvc.Metadata.Namespace || pod.NodeName == "" || !podReferencesPVC(pod, pvc.Metadata.Name) {
393+
continue
394+
}
395+
addNode(pod.NodeName)
396+
}
397+
}
398+
out := make([]string, 0, len(seen))
399+
for nodeName := range seen {
400+
out = append(out, nodeName)
401+
}
402+
sort.Strings(out)
403+
return out
404+
}
405+
406+
func podReferencesPVC(pod k8sresources.Pod, pvcName string) bool {
407+
for _, ref := range pod.PVCRefs {
408+
if ref == pvcName {
409+
return true
410+
}
411+
}
412+
return false
413+
}
414+
377415
func roleBindingID(cluster string, roleBinding k8sresources.RoleBinding) model.CanonicalID {
378416
return model.NewCanonicalID(model.ResourceRef{Cluster: cluster, Group: "rbac.authorization.k8s.io", Kind: "RoleBinding", Namespace: roleBinding.Metadata.Namespace, Name: roleBinding.Metadata.Name, UID: roleBinding.Metadata.UID})
379417
}

internal/graph/pvc_contract_test.go

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,70 @@ func TestBuilderUsesConfiguredCSIComponentRule(t *testing.T) {
146146
}
147147
}
148148

149+
func TestBuilderScopesPVCSINodeAgentToConsumingPodNode(t *testing.T) {
150+
builder := graph.NewBuilder("cluster-a")
151+
builder.SetCSIComponentRules(openLocalCSIComponentRules())
152+
snapshot := collectk8s.Snapshot{
153+
Pods: []resources.Pod{
154+
{
155+
Metadata: resources.Metadata{UID: "pod-uid", Name: "app-0", Namespace: "default"},
156+
NodeName: "node-a",
157+
PVCRefs: []string{"data"},
158+
},
159+
{
160+
Metadata: resources.Metadata{UID: "csi-agent-a-uid", Name: "open-local-agent-node-a", Namespace: "kube-system"},
161+
NodeName: "node-a",
162+
},
163+
{
164+
Metadata: resources.Metadata{UID: "csi-agent-b-uid", Name: "open-local-agent-node-b", Namespace: "kube-system"},
165+
NodeName: "node-b",
166+
},
167+
},
168+
PVCs: []resources.PVC{{
169+
Metadata: resources.Metadata{UID: "pvc-uid", Name: "data", Namespace: "default"},
170+
VolumeName: "pv-data",
171+
StorageClassName: "open-local",
172+
Status: "Bound",
173+
}},
174+
PVs: []resources.PV{{
175+
Metadata: resources.Metadata{UID: "pv-uid", Name: "pv-data"},
176+
StorageClassName: "open-local",
177+
Status: "Bound",
178+
CSI: map[string]string{"driver": "local.csi.aliyun.com", "handle": "vol-123"},
179+
}},
180+
StorageClasses: []resources.StorageClass{{
181+
Metadata: resources.Metadata{UID: "sc-uid", Name: "open-local"},
182+
Provisioner: "local.csi.aliyun.com",
183+
}},
184+
}
185+
186+
nodes, edges := builder.Build(snapshot)
187+
agentA, ok := nodeIDByName(nodes, model.NodeKindPod, "kube-system", "open-local-agent-node-a")
188+
if !ok {
189+
t.Fatal("expected node-a CSI agent")
190+
}
191+
agentB, ok := nodeIDByName(nodes, model.NodeKindPod, "kube-system", "open-local-agent-node-b")
192+
if !ok {
193+
t.Fatal("expected node-b CSI agent")
194+
}
195+
servedByEdges := 0
196+
for _, edge := range edges {
197+
if edge.Kind != model.EdgeKindServedByCSINodeAgent {
198+
continue
199+
}
200+
servedByEdges++
201+
if edge.To == agentB {
202+
t.Fatal("did not expect PV to be served by off-node CSI agent")
203+
}
204+
if edge.To != agentA {
205+
t.Fatalf("expected PV to be served by node-a CSI agent, got %s", edge.To)
206+
}
207+
}
208+
if servedByEdges != 1 {
209+
t.Fatalf("expected one scoped PV node-agent edge, got %d", servedByEdges)
210+
}
211+
}
212+
149213
func storageTopologySnapshot() collectk8s.Snapshot {
150214
return collectk8s.Snapshot{
151215
Pods: []resources.Pod{
@@ -192,6 +256,15 @@ func containsNodeKind(nodes []model.Node, kind model.NodeKind) bool {
192256
return false
193257
}
194258

259+
func nodeIDByName(nodes []model.Node, kind model.NodeKind, namespace, name string) (model.CanonicalID, bool) {
260+
for _, node := range nodes {
261+
if node.Kind == kind && node.Namespace == namespace && node.Name == name {
262+
return node.ID, true
263+
}
264+
}
265+
return "", false
266+
}
267+
195268
func openLocalCSIComponentRules() []infer.CSIComponentRule {
196269
return []infer.CSIComponentRule{{
197270
Driver: "local.csi.aliyun.com",

internal/model/relation_spec.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -286,7 +286,7 @@ var relationSpecs = []RelationSpec{
286286
},
287287
{
288288
Kind: EdgeKindServedByCSINodeAgent,
289-
Comment: "Relates a PersistentVolume to the CSI node-agent Pod on its affinity node.",
289+
Comment: "Relates a PersistentVolume to a CSI node-agent Pod on a PV affinity node or consuming Pod node.",
290290
Domain: string(NodeKindPV),
291291
Range: string(NodeKindPod),
292292
DefaultSourceType: EdgeSourceTypeInferenceRule,

internal/reconcile/storage.go

Lines changed: 41 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@ package reconcile
22

33
import (
44
"fmt"
5+
"sort"
6+
"strings"
57

68
collectk8s "github.com/Colvin-Y/kubernetes-ontology/internal/collect/k8s"
79
"github.com/Colvin-Y/kubernetes-ontology/internal/collect/k8s/resources"
@@ -235,7 +237,7 @@ func (r *StorageReconciler) rebuildStorageEdges(snapshot collectk8s.Snapshot, cu
235237
}
236238
pvNode := pvNode(r.cluster, pv)
237239
pvNode.ID = pvID
238-
correlation := correlator.Correlate(pvNode, pv.CSI["nodeAffinity"], infraPods)
240+
correlation := infer.CorrelatePVToCSIComponents(correlator, pvNode, csiAgentNodeNamesForPV(pv, snapshot), infraPods)
239241
for _, edge := range correlation.Edges {
240242
if err := r.kernel.UpsertEdge(edge); err != nil {
241243
return upserted, err
@@ -409,6 +411,44 @@ func storageClassNameForPVC(pvc resources.PVC, pvs []resources.PV) string {
409411
return ""
410412
}
411413

414+
func csiAgentNodeNamesForPV(pv resources.PV, snapshot collectk8s.Snapshot) []string {
415+
seen := make(map[string]struct{})
416+
addNode := func(nodeName string) {
417+
nodeName = strings.TrimSpace(nodeName)
418+
if nodeName == "" {
419+
return
420+
}
421+
seen[nodeName] = struct{}{}
422+
}
423+
addNode(pv.CSI["nodeAffinity"])
424+
for _, pvc := range snapshot.PVCs {
425+
if pvc.VolumeName != pv.Metadata.Name {
426+
continue
427+
}
428+
for _, pod := range snapshot.Pods {
429+
if pod.Metadata.Namespace != pvc.Metadata.Namespace || pod.NodeName == "" || !podReferencesPVC(pod, pvc.Metadata.Name) {
430+
continue
431+
}
432+
addNode(pod.NodeName)
433+
}
434+
}
435+
out := make([]string, 0, len(seen))
436+
for nodeName := range seen {
437+
out = append(out, nodeName)
438+
}
439+
sort.Strings(out)
440+
return out
441+
}
442+
443+
func podReferencesPVC(pod resources.Pod, pvcName string) bool {
444+
for _, ref := range pod.PVCRefs {
445+
if ref == pvcName {
446+
return true
447+
}
448+
}
449+
return false
450+
}
451+
412452
func storageInfraPods(cluster string, snapshot collectk8s.Snapshot) []model.Node {
413453
out := make([]model.Node, 0)
414454
for _, pod := range snapshot.Pods {

internal/resolve/infer/csi.go

Lines changed: 49 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package infer
22

33
import (
44
"fmt"
5+
"sort"
56
"strings"
67

78
"github.com/Colvin-Y/kubernetes-ontology/internal/model"
@@ -106,6 +107,35 @@ func NewCSIComponentRegistry(rules []CSIComponentRule) *Registry {
106107
return NewRegistry(correlators...)
107108
}
108109

110+
func CorrelatePVToCSIComponents(correlator CSICorrelator, pv model.Node, nodeNames []string, infraPods []model.Node) CorrelationResult {
111+
normalizedNodeNames := normalizeNodeNames(nodeNames)
112+
if len(normalizedNodeNames) == 0 {
113+
return correlator.Correlate(pv, "", infraPods)
114+
}
115+
116+
result := CorrelationResult{Edges: make([]model.Edge, 0), Evidence: make([]string, 0)}
117+
seenEdges := make(map[string]struct{})
118+
seenEvidence := make(map[string]struct{})
119+
for _, nodeName := range normalizedNodeNames {
120+
correlation := correlator.Correlate(pv, nodeName, infraPods)
121+
for _, edge := range correlation.Edges {
122+
if _, seen := seenEdges[edge.Key()]; seen {
123+
continue
124+
}
125+
result.Edges = append(result.Edges, edge)
126+
seenEdges[edge.Key()] = struct{}{}
127+
}
128+
for _, evidence := range correlation.Evidence {
129+
if _, seen := seenEvidence[evidence]; seen {
130+
continue
131+
}
132+
result.Evidence = append(result.Evidence, evidence)
133+
seenEvidence[evidence] = struct{}{}
134+
}
135+
}
136+
return result
137+
}
138+
109139
func IsCSIProvisioner(provisioner string, observedCSIDriver bool, rules []CSIComponentRule) bool {
110140
if provisioner == "" {
111141
return false
@@ -187,9 +217,9 @@ func (c componentRuleCorrelator) Correlate(pv model.Node, affinityNodeName strin
187217
}
188218
if len(c.rule.NodeAgentPodPrefixes) > 0 {
189219
if affinityNodeName == "" {
190-
result.Evidence = append(result.Evidence, fmt.Sprintf("csi: PV affinity node missing for driver %s", c.rule.Driver))
220+
result.Evidence = append(result.Evidence, fmt.Sprintf("csi: PV affinity or consuming pod node missing for driver %s", c.rule.Driver))
191221
} else if !foundAgent {
192-
result.Evidence = append(result.Evidence, fmt.Sprintf("csi: no node agent found for driver %s on PV affinity node %s", c.rule.Driver, affinityNodeName))
222+
result.Evidence = append(result.Evidence, fmt.Sprintf("csi: no node agent found for driver %s on node %s", c.rule.Driver, affinityNodeName))
193223
}
194224
}
195225
return result
@@ -217,6 +247,23 @@ func hasAnyPrefix(value string, prefixes []string) bool {
217247
return false
218248
}
219249

250+
func normalizeNodeNames(nodeNames []string) []string {
251+
seen := make(map[string]struct{}, len(nodeNames))
252+
for _, nodeName := range nodeNames {
253+
nodeName = strings.TrimSpace(nodeName)
254+
if nodeName == "" {
255+
continue
256+
}
257+
seen[nodeName] = struct{}{}
258+
}
259+
out := make([]string, 0, len(seen))
260+
for nodeName := range seen {
261+
out = append(out, nodeName)
262+
}
263+
sort.Strings(out)
264+
return out
265+
}
266+
220267
func looksLikeCSIProvisioner(provisioner string) bool {
221268
if strings.HasPrefix(provisioner, "kubernetes.io/") {
222269
return false

internal/resolve/infer/csi_test.go

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -121,18 +121,35 @@ func TestCSIComponentRuleCorrelatesPVToComponents(t *testing.T) {
121121
Namespace: "storage-system",
122122
Attributes: map[string]any{"nodeName": "node-a"},
123123
}
124+
otherAgent := model.Node{
125+
ID: model.NewCanonicalID(model.ResourceRef{Cluster: "cluster-a", Group: "core", Kind: "Pod", Namespace: "storage-system", Name: "disk-agent-node-b"}),
126+
Kind: model.NodeKindPod,
127+
Name: "disk-agent-node-b",
128+
Namespace: "storage-system",
129+
Attributes: map[string]any{"nodeName": "node-b"},
130+
}
124131

125-
result := correlator.Correlate(pv, "node-a", []model.Node{controller, agent})
132+
result := correlator.Correlate(pv, "node-a", []model.Node{controller, agent, otherAgent})
126133
kinds := map[model.EdgeKind]bool{}
134+
servedByAgentEdges := 0
127135
for _, edge := range result.Edges {
128136
kinds[edge.Kind] = true
137+
if edge.Kind == model.EdgeKindServedByCSINodeAgent {
138+
servedByAgentEdges++
139+
if edge.To != agent.ID {
140+
t.Fatalf("expected PV node-agent edge to target same-node agent, got %s", edge.To)
141+
}
142+
}
129143
}
130144
if !kinds[model.EdgeKindManagedByCSIController] {
131145
t.Fatal("expected PV controller edge")
132146
}
133147
if !kinds[model.EdgeKindServedByCSINodeAgent] {
134148
t.Fatal("expected PV node-agent edge")
135149
}
150+
if servedByAgentEdges != 1 {
151+
t.Fatalf("expected exactly one PV node-agent edge, got %d", servedByAgentEdges)
152+
}
136153
}
137154

138155
func TestParseCSIComponentRules(t *testing.T) {

0 commit comments

Comments
 (0)