From 135782e429fda199ecd69327648ad1ff643c23ff Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Wed, 24 Sep 2025 18:02:44 +0800 Subject: [PATCH 1/9] adjust directref relation update logic need to split to atomic reference pair to do sort and compare --- resourcetopo/node_topology_updater.go | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/resourcetopo/node_topology_updater.go b/resourcetopo/node_topology_updater.go index d4a37f2..947c777 100644 --- a/resourcetopo/node_topology_updater.go +++ b/resourcetopo/node_topology_updater.go @@ -289,3 +289,21 @@ func (s *nodeStorage) removeLabelResourceRelation(node *nodeInfo, relation *Reso } } } + +func (s *nodeStorage) removeDirectResourceRelation(node *nodeInfo, relation *ResourceRelation) { + postStorage := s.manager.getStorage(relation.PostMeta) + if postStorage == nil { + klog.Error("Failed to get node Storage by %s, ignore this delete request", + generateMetaKey(relation.PostMeta)) + return + } + if len(relation.DirectRefs) > 0 { + for _, ref := range relation.DirectRefs { + postNode := postStorage.getNode(relation.Cluster, ref.Namespace, ref.Name) + if deleteDirectRelation(node, postNode) { + node.noticePostOrderRelationDeleted(postNode) + postNode.checkGC() + } + } + } +} From 67bb35a1d62f2561419357c4e7338607c3e2d817 Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Wed, 24 Sep 2025 18:17:11 +0800 Subject: [PATCH 2/9] rm unused func --- resourcetopo/node_topology_updater.go | 18 ------------------ 1 file changed, 18 deletions(-) diff --git a/resourcetopo/node_topology_updater.go b/resourcetopo/node_topology_updater.go index 947c777..d4a37f2 100644 --- a/resourcetopo/node_topology_updater.go +++ b/resourcetopo/node_topology_updater.go @@ -289,21 +289,3 @@ func (s *nodeStorage) removeLabelResourceRelation(node *nodeInfo, relation *Reso } } } - -func (s *nodeStorage) removeDirectResourceRelation(node *nodeInfo, relation *ResourceRelation) { - postStorage := s.manager.getStorage(relation.PostMeta) - if postStorage == nil { - klog.Error("Failed to get node Storage by %s, ignore this delete request", - generateMetaKey(relation.PostMeta)) - return - } - if len(relation.DirectRefs) > 0 { - for _, ref := range relation.DirectRefs { - postNode := postStorage.getNode(relation.Cluster, ref.Namespace, ref.Name) - if deleteDirectRelation(node, postNode) { - node.noticePostOrderRelationDeleted(postNode) - postNode.checkGC() - } - } - } -} From ba4939665f36b6cdd3f45f7c60177b0b309e9cbb Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Mon, 16 Mar 2026 17:29:12 +0800 Subject: [PATCH 3/9] refactor(resourcetopo): fix lock hierarchy and add event handler timeout This commit addresses potential deadlock issues and adds timeout protection for event handlers. Lock Hierarchy Fixes: - Add doc.go with comprehensive lock hierarchy documentation (7 levels) - Fix checkGC(): acquire lock before metaLock, release between - Fix checkLabelUpdateForPostNode(): copy data under lock, release before processing - Fix updateNodeMeta(): release metaLock before acquiring lock - Fix deleteAllRelation(): restructure to follow preNode->postNode lock order - Add storagesLock RWMutex to protect storages map from concurrent access Event Handler Timeout: - Add EventHandlerTimeout config option (default 1s) - Implement timeout in handleNodeEvent() with context cancellation - Implement timeout in handleRelationEvent() with context cancellation - Handlers run sequentially; timeout skips remaining handlers after warning Other Changes: - Change event queue log level from Infof to V(6).Infof - Add threading model documentation Co-Authored-By: Claude Opus 4.6 --- resourcetopo/doc.go | 95 +++++++++++++++++++++ resourcetopo/event_queue.go | 123 ++++++++++++++++++++++++++-- resourcetopo/manager.go | 54 +++++++++--- resourcetopo/node_info.go | 43 +++++++--- resourcetopo/node_relation_utils.go | 82 +++++++++++-------- resourcetopo/types.go | 6 ++ 6 files changed, 338 insertions(+), 65 deletions(-) create mode 100644 resourcetopo/doc.go diff --git a/resourcetopo/doc.go b/resourcetopo/doc.go new file mode 100644 index 0000000..cb3895c --- /dev/null +++ b/resourcetopo/doc.go @@ -0,0 +1,95 @@ +/** + * Copyright 2024 KusionStack Authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +// Package resourcetopo provides a resource topology management system for Kubernetes-like +// resources. It maintains relationships between different resource types and notifies +// handlers when resources or their relationships change. +// +// Lock Hierarchy +// +// To prevent deadlocks, locks must be acquired in the following order: +// +// Level 1: manager.configLock +// - Protects: started flag +// - Used during: topology configuration, Start() +// +// Level 2: manager.storagesLock +// - Protects: storages map +// - Lock type: RWMutex (RLock for reads, Lock for writes) +// +// Level 3: nodeStorage.storageLock +// - Protects: clusterNodes map +// - Note: Never acquire while holding any nodeInfo.lock +// +// Level 4: nodeInfo.lock +// - Protects: PreOrders/PostOrders relation lists (directReferredPreOrders, etc.) +// - Rule: Always acquire preNode.lock before postNode.lock +// +// Level 5: nodeInfo.relationsLock +// - Protects: labelRelations +// - Rule: Acquire before metaLock if both needed +// +// Level 6: nodeInfo.metaLock +// - Protects: labels, ownerNodes, objectExisted +// - This is the innermost lock - never acquire other locks while holding metaLock +// +// Level 7: nodeStorage.handlersLock +// - Protects: nodeUpdateHandler, relationUpdateHandler +// - Independent - can be acquired anytime +// +// Threading Model +// +// This package uses multiple goroutines for concurrent processing. Understanding the +// threading model is essential for avoiding race conditions and deadlocks. +// +// 1. Informer Threads (External): +// - Source: Kubernetes informers (one per resource type) +// - Entry points: nodeStorage.OnAdd, OnUpdate, OnDelete +// - Purpose: Receives resource change events from API server +// - Lock behavior: Acquires metaLock.Lock() during updateNodeMeta(), relationsLock during relation updates +// - Note: These are the ONLY threads that write to metaLock (labels, ownerRefs, objectExisted) +// +// 2. Event Processor Threads (Internal): +// - Source: Created by Start() -> startHandleEvent() +// - Count: 2 goroutines +// - handleNodeEvent(): Processes node add/update/delete/relatedUpdate events +// - handleRelationEvent(): Processes relation add/delete events +// - Trigger: newNodeEvent(), newRelationEvent() queue events via workqueue +// - Lock behavior: Read-only access to handlers via RLock, no direct node locking +// - Note: Reads handlers while holding no locks - relies on handlersLock for registration +// +// 3. User/Client Threads (External): +// - Source: User code calling Manager methods +// - Examples: GetNode(), GetTopoNodeStorage(), AddNodeHandler() +// - Lock behavior: Uses storagesLock.RLock() for reads +// - Note: These are typically short-lived operations +// +// 4. Callback Threads (External - via Event Processors): +// - Source: User-provided NodeHandler and RelationHandler callbacks +// - Entry: Called by handleNodeEvent/handleRelationEvent +// - Lock behavior: No locks held during callback execution +// - Safety: Handlers receive node references but should not cache them long-term +// +// Thread Safety Guidelines +// +// - NodeInfo references returned by GetNode() are safe for concurrent reads +// - Do NOT call AddTopologyConfig() after Start() - not thread safe +// - Handler registration (AddNodeHandler, AddRelationHandler) is thread safe +// - All callbacks (OnAdd, OnUpdate, etc.) are called synchronously - do not block +// - nodeInfo.lock can be held for extended periods during relation changes +// - metaLock is never held during callbacks or cross-node operations +// +package resourcetopo diff --git a/resourcetopo/event_queue.go b/resourcetopo/event_queue.go index 8d653d6..b06fbb4 100644 --- a/resourcetopo/event_queue.go +++ b/resourcetopo/event_queue.go @@ -17,6 +17,7 @@ package resourcetopo import ( + "context" "strings" "time" @@ -41,6 +42,7 @@ const ( defaultRelationEventHandleRateMaxDelay = time.Second defaultNodeEventHandlePeriod = time.Second defaultRelationEventHandlePeriod = time.Second + defaultEventHandlerTimeout = time.Second ) func (m *manager) startHandleEvent(stopCh <-chan struct{}) { @@ -51,43 +53,108 @@ func (m *manager) startHandleEvent(stopCh <-chan struct{}) { func (m *manager) handleNodeEvent() { for { item, shutdown := m.nodeEventQueue.Get() + klog.V(6).Infof("handleNodeEvent: item: %v", item) if shutdown { return } info, ok := item.(string) if !ok { klog.Errorf("Unexpected node event queue item %v", item) + m.nodeEventQueue.Done(item) continue } evtType, node := m.decodeString2NodeEvent(info) if node == nil { + m.nodeEventQueue.Done(item) continue } storage := node.storageRef if storage == nil { klog.Errorf("Unexpected nil nodeStorage for nodeEvent node %v", node) + m.nodeEventQueue.Done(item) continue } + // Create context with timeout for all handlers in this event + ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) + timeout := false + switch evtType { case EventTypeAdd: for _, h := range storage.nodeUpdateHandler { - h.OnAdd(node) + if timeout { + break + } + done := make(chan struct{}) + go func(handler NodeHandler) { + handler.OnAdd(node) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + klog.Warningf("Node handler OnAdd timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + } } case EventTypeUpdate: for _, h := range storage.nodeUpdateHandler { - h.OnUpdate(node) + if timeout { + break + } + done := make(chan struct{}) + go func(handler NodeHandler) { + handler.OnUpdate(node) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + klog.Warningf("Node handler OnUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + } } case EventTypeDelete: for _, h := range storage.nodeUpdateHandler { - h.OnDelete(node) + if timeout { + break + } + done := make(chan struct{}) + go func(handler NodeHandler) { + handler.OnDelete(node) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + klog.Warningf("Node handler OnDelete timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + } } case EventTypeRelatedUpdate: for _, h := range storage.nodeUpdateHandler { - h.OnRelatedUpdate(node) + if timeout { + break + } + done := make(chan struct{}) + go func(handler NodeHandler) { + handler.OnRelatedUpdate(node) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + klog.Warningf("Node handler OnRelatedUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + } } } + cancel() m.nodeEventQueue.Done(item) } } @@ -102,39 +169,79 @@ func (m *manager) handleRelationEvent() { info, ok := item.(string) if !ok { klog.Errorf("Unexpected relation event queue item %v", item) + m.relationEventQueue.Done(item) continue } evtType, preNode, postNode := m.decodeString2RelationEvent(info) if preNode == nil || postNode == nil { + m.relationEventQueue.Done(item) continue } storage := preNode.storageRef if storage == nil { klog.Errorf("Unexpected nil nodeStorage for relaltion event preNode node %v", preNode) + m.relationEventQueue.Done(item) continue } handlers := storage.relationUpdateHandler[postNode.storageRef.metaKey] if handlers == nil { + m.relationEventQueue.Done(item) continue } + // Create context with timeout for all handlers in this event + ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) + timeout := false + switch evtType { case EventTypeAdd: - for _, handler := range handlers { - handler.OnAdd(preNode, postNode) + for _, h := range handlers { + if timeout { + break + } + done := make(chan struct{}) + go func(handler RelationHandler) { + handler.OnAdd(preNode, postNode) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + klog.Warningf("Relation handler OnAdd timeout for relation %s/%s -> %s/%s, skipping remaining handlers", + preNode.namespace, preNode.name, postNode.namespace, postNode.name) + } } case EventTypeDelete: - for _, handler := range handlers { - handler.OnDelete(preNode, postNode) + for _, h := range handlers { + if timeout { + break + } + done := make(chan struct{}) + go func(handler RelationHandler) { + handler.OnDelete(preNode, postNode) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + klog.Warningf("Relation handler OnDelete timeout for relation %s/%s -> %s/%s, skipping remaining handlers", + preNode.namespace, preNode.name, postNode.namespace, postNode.name) + } } } + cancel() m.relationEventQueue.Done(item) } } func (m *manager) newNodeEvent(info *nodeInfo, eType eventType) { key := m.encodeNodeEvent2String(eType, info) + klog.V(6).Infof("New node event queue key: %v", key) m.nodeEventQueue.AddRateLimited(key) } diff --git a/resourcetopo/manager.go b/resourcetopo/manager.go index d9ddacc..e17653a 100644 --- a/resourcetopo/manager.go +++ b/resourcetopo/manager.go @@ -20,6 +20,7 @@ import ( "container/list" "fmt" "sync" + "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" @@ -35,7 +36,10 @@ type manager struct { configLock sync.Mutex started bool - storages map[string]*nodeStorage // meta => nodeStorage and status info + storages map[string]*nodeStorage // meta => nodeStorage and status info + storagesLock sync.RWMutex // protects storages map + + eventHandlerTimeout time.Duration // timeout for event handler callbacks } func NewResourcesTopoManager(cfg ManagerConfig) (Manager, error) { @@ -44,9 +48,10 @@ func NewResourcesTopoManager(cfg ManagerConfig) (Manager, error) { nodeRatelimiter := workqueue.NewItemExponentialFailureRateLimiter(cfg.NodeEventHandleRateMinDelay, cfg.NodeEventHandleRateMaxDelay) relationRatelimiter := workqueue.NewItemExponentialFailureRateLimiter(cfg.RelationEventHandleRateMinDelay, cfg.RelationEventHandleRateMaxDelay) m := &manager{ - nodeEventQueue: workqueue.NewNamedRateLimitingQueue(nodeRatelimiter, "resourcetopoNodeEventQueue"), - relationEventQueue: workqueue.NewNamedRateLimitingQueue(relationRatelimiter, "resourcetopoRelaltionEventQueue"), - storages: make(map[string]*nodeStorage), + nodeEventQueue: workqueue.NewNamedRateLimitingQueue(nodeRatelimiter, "resourcetopoNodeEventQueue"), + relationEventQueue: workqueue.NewNamedRateLimitingQueue(relationRatelimiter, "resourcetopoRelaltionEventQueue"), + storages: make(map[string]*nodeStorage), + eventHandlerTimeout: cfg.EventHandlerTimeout, } if cfg.TopologyConfig != nil { @@ -102,7 +107,9 @@ func (m *manager) AddTopologyConfig(cfg TopologyConfig) error { // the two type nodes configured by preMeta and postMeta func (m *manager) AddRelationHandler(preOrder, postOrder metav1.TypeMeta, handler RelationHandler) error { preOrderKey := generateMetaKey(preOrder) + m.storagesLock.RLock() preStorage, ok := m.storages[preOrderKey] + m.storagesLock.RUnlock() if !ok { return fmt.Errorf("failed to find initialized nodeStorage for resource %s", preOrderKey) } @@ -112,7 +119,9 @@ func (m *manager) AddRelationHandler(preOrder, postOrder metav1.TypeMeta, handle // AddNodeHandler add a new node handler for nodes of type configured by meta. func (m *manager) AddNodeHandler(meta metav1.TypeMeta, handler NodeHandler) error { + m.storagesLock.RLock() s := m.storages[generateMetaKey(meta)] + m.storagesLock.RUnlock() if s == nil { return fmt.Errorf("resource %v not configured in this manager", meta) } @@ -135,18 +144,22 @@ func (m *manager) Start(stopCh <-chan struct{}) { // GetTopoNodeStorage return the ref to TopoNodeStorage that match resource meta. func (m *manager) GetTopoNodeStorage(meta metav1.TypeMeta) (TopoNodeStorage, error) { - if s := m.storages[generateMetaKey(meta)]; s == nil { + m.storagesLock.RLock() + s := m.storages[generateMetaKey(meta)] + m.storagesLock.RUnlock() + if s == nil { return nil, fmt.Errorf("resource %v not configured in this manager", meta) - } else { - return s, nil } + return s, nil } // GetNode return the ref to Node that match the resouorce meta and node's name if existed. func (m *manager) GetNode(meta metav1.TypeMeta, name types.NamespacedName) (NodeInfo, error) { - s, err := m.GetTopoNodeStorage(meta) - if err != nil { - return nil, err + m.storagesLock.RLock() + s := m.storages[generateMetaKey(meta)] + m.storagesLock.RUnlock() + if s == nil { + return nil, fmt.Errorf("resource %v not configured in this manager", meta) } return s.GetNode(name) } @@ -157,7 +170,9 @@ func (m *manager) getOrCreateStorage(typeMeta metav1.TypeMeta, getInformer func( } key := generateMetaKey(typeMeta) + m.storagesLock.RLock() s, ok := m.storages[key] + m.storagesLock.RUnlock() if ok { return s, nil } @@ -168,25 +183,32 @@ func (m *manager) getOrCreateStorage(typeMeta metav1.TypeMeta, getInformer func( } s = newNodeStorage(m, informer, typeMeta) + m.storagesLock.Lock() m.storages[key] = s + m.storagesLock.Unlock() return s, nil } func (m *manager) getStorage(meta metav1.TypeMeta) *nodeStorage { key := generateMetaKey(meta) - + m.storagesLock.RLock() + defer m.storagesLock.RUnlock() return m.storages[key] } func (m *manager) createVirtualStorage(meta metav1.TypeMeta) (*nodeStorage, error) { key := generateMetaKey(meta) + m.storagesLock.RLock() _, ok := m.storages[key] + m.storagesLock.RUnlock() if ok { return nil, fmt.Errorf("unexpected twice configuration for virtual resource %s", key) } s := newVirtualStorage(m, meta) + m.storagesLock.Lock() m.storages[key] = s + m.storagesLock.Unlock() return s, nil } @@ -195,7 +217,14 @@ func (m *manager) dagCheck() error { visited := make(map[string]bool) stack := list.New() + m.storagesLock.RLock() + storagesCopy := make(map[string]*nodeStorage, len(m.storages)) for k, s := range m.storages { + storagesCopy[k] = s + } + m.storagesLock.RUnlock() + + for k, s := range storagesCopy { if visited[k] { if existInList(stack, k) { return fmt.Errorf("DAG check for resource %s failed", k) @@ -223,4 +252,7 @@ func checkManagerConfig(c *ManagerConfig) { if c.RelationEventHandleRateMaxDelay < c.RelationEventHandleRateMinDelay { c.RelationEventHandleRateMaxDelay = defaultRelationEventHandleRateMaxDelay } + if c.EventHandlerTimeout <= 0 { + c.EventHandlerTimeout = defaultEventHandlerTimeout + } } diff --git a/resourcetopo/node_info.go b/resourcetopo/node_info.go index 4e40c43..23364d9 100644 --- a/resourcetopo/node_info.go +++ b/resourcetopo/node_info.go @@ -129,12 +129,17 @@ func (node *nodeInfo) checkLabelUpdateForPostNode(postNode *nodeInfo) { return } + // Copy labelRelations under lock to avoid holding lock during relation operations node.relationsLock.RLock() - defer node.relationsLock.RUnlock() if len(node.labelRelations) == 0 { + node.relationsLock.RUnlock() return } - for _, relation := range node.labelRelations { + labelRelationsCopy := make([]ResourceRelation, len(node.labelRelations)) + copy(labelRelationsCopy, node.labelRelations) + node.relationsLock.RUnlock() + + for _, relation := range labelRelationsCopy { if !typeEqual(relation.PostMeta, postNode.storageRef.meta) { return } @@ -157,7 +162,6 @@ func (node *nodeInfo) checkLabelUpdateForPostNode(postNode *nodeInfo) { func (n *nodeInfo) updateNodeMeta(obj Object) { n.metaLock.Lock() - defer n.metaLock.Unlock() n.labels = obj.GetLabels() var ownerInfos []ownerInfo @@ -170,8 +174,13 @@ func (n *nodeInfo) updateNodeMeta(obj Object) { } n.ownerNodes = ownerInfos - if !n.objectExisted { - n.objectExisted = true + objectJustCreated := !n.objectExisted + n.objectExisted = true + n.metaLock.Unlock() + + // Acquire lock for relation iteration after releasing metaLock + // to maintain lock ordering: lock before metaLock + if objectJustCreated { n.lock.RLock() defer n.lock.RUnlock() rangeNodeList(n.directReferredPreOrders, func(preOrder *nodeInfo) { @@ -236,17 +245,25 @@ func (n *nodeInfo) ownersEqualed(reference []metav1.OwnerReference) bool { } func (n *nodeInfo) checkGC() { - n.metaLock.RLock() - defer n.metaLock.RUnlock() + // Acquire lock for relation lists first (Level 4) n.lock.RLock() - defer n.lock.RUnlock() - - if !n.objectExisted && - // n maybe set as notExisted, but still have related nodes to handle - (n.directReferredPreOrders == nil || n.directReferredPreOrders.Len() == 0) && + relationsEmpty := (n.directReferredPreOrders == nil || n.directReferredPreOrders.Len() == 0) && (n.directReferredPostOrders == nil || n.directReferredPostOrders.Len() == 0) && (n.labelReferredPreOrders == nil || n.labelReferredPreOrders.Len() == 0) && - (n.labelReferredPostOrders == nil || n.labelReferredPostOrders.Len() == 0) { + (n.labelReferredPostOrders == nil || n.labelReferredPostOrders.Len() == 0) + n.lock.RUnlock() + + if !relationsEmpty { + return + } + + // Then acquire metaLock to check object existence (Level 6) + n.metaLock.Lock() + objectExisted := n.objectExisted + n.metaLock.Unlock() + + // n maybe set as notExisted, but still have related nodes to handle + if !objectExisted { n.storageRef.deleteNode(n.cluster, n.namespace, n.name) } } diff --git a/resourcetopo/node_relation_utils.go b/resourcetopo/node_relation_utils.go index 75a1221..ec3fc37 100644 --- a/resourcetopo/node_relation_utils.go +++ b/resourcetopo/node_relation_utils.go @@ -95,61 +95,77 @@ func deleteDirectRelation(preNode, postNode *nodeInfo) bool { // deleteAllRelation will remove all relation in this node. // only be called in object deletion scenario. func deleteAllRelation(node *nodeInfo) { - // In resourcetopo package, when do relation change operations, - // we always lock preOrder before postOrder, expect in this function. - // but to keep transition atomic, and in object deletion scenario, - // this is okay to hold the lock in whole process, and lock preOrderNode as needed. - // node.metaLock.Lock() - // defer node.metaLock.Unlock() - node.relationsLock.Lock() + // First, mark object as deleted under metaLock (Level 6 - innermost) node.objectDeleted() - node.labelRelations = nil - defer node.relationsLock.Unlock() - - node.lock.Lock() - defer node.lock.Unlock() + // Collect all related nodes before acquiring locks to avoid holding locks during iteration + var directPostNodes, labelPostNodes, labelPreNodes, directPreNodes []*nodeInfo + node.lock.RLock() rangeNodeList(node.directReferredPostOrders, func(postNode *nodeInfo) { - postNode.lock.Lock() - removeFromList(postNode.directReferredPreOrders, node) - postNode.lock.Unlock() + directPostNodes = append(directPostNodes, postNode) + }) + rangeNodeList(node.labelReferredPostOrders, func(postNode *nodeInfo) { + labelPostNodes = append(labelPostNodes, postNode) + }) + rangeNodeList(node.labelReferredPreOrders, func(preNode *nodeInfo) { + labelPreNodes = append(labelPreNodes, preNode) + }) + rangeNodeList(node.directReferredPreOrders, func(preNode *nodeInfo) { + directPreNodes = append(directPreNodes, preNode) + }) + node.lock.RUnlock() + // Process post-nodes (node is preOrder, so lock node first, then postNode) + // This follows the rule: preNode.lock before postNode.lock + for _, postNode := range directPostNodes { + deleteDirectRelation(node, postNode) node.noticePostOrderRelationDeleted(postNode) postNode.checkGC() - }) - node.directReferredPostOrders = nil - - rangeNodeList(node.labelReferredPostOrders, func(postNode *nodeInfo) { - postNode.lock.Lock() - removeFromList(postNode.labelReferredPreOrders, node) - postNode.lock.Unlock() + } + for _, postNode := range labelPostNodes { + deleteLabelRelation(node, postNode) node.noticePostOrderRelationDeleted(postNode) - }) - node.labelReferredPostOrders = nil + } - // node has been set as deleted, labels set to nil, - // so can not be locked as postOrder to add new label selector relation - rangeNodeList(node.labelReferredPreOrders, func(preNode *nodeInfo) { + // Process pre-nodes (preNode is preOrder, so lock preNode first, then node) + // This follows the rule: preNode.lock before postNode.lock + for _, preNode := range labelPreNodes { + // For label relations, node is postOrder, so we need preNode.lock first preNode.lock.Lock() removeFromList(preNode.labelReferredPostOrders, node) preNode.lock.Unlock() node.noticePreOrderRelationDeleted(preNode) - }) - node.labelReferredPreOrders = nil + } - rangeNodeList(node.directReferredPreOrders, func(preNode *nodeInfo) { + for _, preNode := range directPreNodes { if preNode.storageRef.virtualResource { - // virtual node will not start a relation pair change, so virtual preOrder node can be locked later + // For virtual resources, we can modify the relation directly preNode.lock.Lock() removeFromList(preNode.directReferredPostOrders, node) preNode.lock.Unlock() + // Also update node's list while holding node's lock + node.lock.Lock() removeFromList(node.directReferredPreOrders, preNode) + node.lock.Unlock() node.noticePreOrderRelationDeleted(preNode) preNode.checkVirtualNodeGC() } else { - // notice relation deleted, but hold this place in case this object will be recreated + // For non-virtual resources, just notify but keep node's list as-is + // in case this object will be recreated node.noticePreOrderRelationDeleted(preNode) } - }) + } + + // Clean up node's relation lists and labelRelations under proper locks + node.relationsLock.Lock() + node.labelRelations = nil + node.relationsLock.Unlock() + + node.lock.Lock() + node.directReferredPostOrders = nil + node.labelReferredPostOrders = nil + node.labelReferredPreOrders = nil + // Keep directReferredPreOrders for non-virtual preNodes in case of recreation + node.lock.Unlock() } diff --git a/resourcetopo/types.go b/resourcetopo/types.go index af8741f..91f42ef 100644 --- a/resourcetopo/types.go +++ b/resourcetopo/types.go @@ -49,6 +49,12 @@ type ManagerConfig struct { NodeEventHandleRateMaxDelay time.Duration RelationEventHandleRateMinDelay time.Duration RelationEventHandleRateMaxDelay time.Duration + + // EventHandlerTimeout is the timeout for each event handler callback. + // If a handler takes longer than this timeout, it will be cancelled and + // remaining handlers will be skipped. A warning will be logged. + // Default: 1 second + EventHandlerTimeout time.Duration } type Manager interface { From 389b6d0c7a03de6152c2bbde5a5fe9f5f21dc4ed Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Mon, 16 Mar 2026 19:24:39 +0800 Subject: [PATCH 4/9] refactor(resourcetopo): use two-value map lookup to prevent NPE Change map get operations to use the two-value lookup pattern (val, ok := map[key]) to explicitly check for key existence. This prevents returning nil values that could cause NPE when callers try to use them. Changes: - manager.go: AddNodeHandler, GetTopoNodeStorage, GetNode now check ok before using value - node_storage.go: getNode and deleteNode check ok for cluster, namespace, and name lookups - GetClusterNode: return nil early if node not found Co-Authored-By: Claude Opus 4.6 --- resourcetopo/manager.go | 12 ++++++------ resourcetopo/node_storage.go | 26 ++++++++++++++++---------- 2 files changed, 22 insertions(+), 16 deletions(-) diff --git a/resourcetopo/manager.go b/resourcetopo/manager.go index e17653a..578e6c5 100644 --- a/resourcetopo/manager.go +++ b/resourcetopo/manager.go @@ -120,9 +120,9 @@ func (m *manager) AddRelationHandler(preOrder, postOrder metav1.TypeMeta, handle // AddNodeHandler add a new node handler for nodes of type configured by meta. func (m *manager) AddNodeHandler(meta metav1.TypeMeta, handler NodeHandler) error { m.storagesLock.RLock() - s := m.storages[generateMetaKey(meta)] + s, ok := m.storages[generateMetaKey(meta)] m.storagesLock.RUnlock() - if s == nil { + if !ok || s == nil { return fmt.Errorf("resource %v not configured in this manager", meta) } s.addNodeHandler(handler) @@ -145,9 +145,9 @@ func (m *manager) Start(stopCh <-chan struct{}) { // GetTopoNodeStorage return the ref to TopoNodeStorage that match resource meta. func (m *manager) GetTopoNodeStorage(meta metav1.TypeMeta) (TopoNodeStorage, error) { m.storagesLock.RLock() - s := m.storages[generateMetaKey(meta)] + s, ok := m.storages[generateMetaKey(meta)] m.storagesLock.RUnlock() - if s == nil { + if !ok || s == nil { return nil, fmt.Errorf("resource %v not configured in this manager", meta) } return s, nil @@ -156,9 +156,9 @@ func (m *manager) GetTopoNodeStorage(meta metav1.TypeMeta) (TopoNodeStorage, err // GetNode return the ref to Node that match the resouorce meta and node's name if existed. func (m *manager) GetNode(meta metav1.TypeMeta, name types.NamespacedName) (NodeInfo, error) { m.storagesLock.RLock() - s := m.storages[generateMetaKey(meta)] + s, ok := m.storages[generateMetaKey(meta)] m.storagesLock.RUnlock() - if s == nil { + if !ok || s == nil { return nil, fmt.Errorf("resource %v not configured in this manager", meta) } return s.GetNode(name) diff --git a/resourcetopo/node_storage.go b/resourcetopo/node_storage.go index a153436..3823e61 100644 --- a/resourcetopo/node_storage.go +++ b/resourcetopo/node_storage.go @@ -90,7 +90,9 @@ func (s *nodeStorage) GetNode(namespacedName types.NamespacedName) (NodeInfo, er func (s *nodeStorage) GetClusterNode(cluster string, namespacedName types.NamespacedName) (NodeInfo, error) { node := s.getNode(cluster, namespacedName.Namespace, namespacedName.Name) - if node != nil && !node.isObjectExisted() { + if node == nil || !node.isObjectExisted() { + // return nil rather than node ref here, since the node may be a virtual node + // todo add a unittest return nil, nil } return node, nil @@ -166,17 +168,21 @@ func (s *nodeStorage) getNode(cluster, namespace, name string) *nodeInfo { s.storageLock.RLock() defer s.storageLock.RUnlock() - clsNodes := s.clusterNodes[cluster] - if clsNodes == nil { + clsNodes, ok := s.clusterNodes[cluster] + if !ok || clsNodes == nil { return nil } - nsNodes := clsNodes.namespaceNodes[namespace] - if nsNodes == nil { + nsNodes, ok := clsNodes.namespaceNodes[namespace] + if !ok || nsNodes == nil { return nil } - return nsNodes[name] + node, ok := nsNodes[name] + if !ok { + return nil + } + return node } func (s *nodeStorage) createNode(cluster, namespace, name string) *nodeInfo { @@ -212,13 +218,13 @@ func (s *nodeStorage) deleteNode(cluster, namespace, name string) { s.storageLock.Lock() defer s.storageLock.Unlock() - clsNodes := s.clusterNodes[cluster] - if clsNodes == nil { + clsNodes, ok := s.clusterNodes[cluster] + if !ok || clsNodes == nil { return } - nsNodes := clsNodes.namespaceNodes[namespace] - if nsNodes == nil { + nsNodes, ok := clsNodes.namespaceNodes[namespace] + if !ok || nsNodes == nil { return } From 67be6435ee3130f4c1e416ce75e451355209203a Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Tue, 17 Mar 2026 14:13:19 +0800 Subject: [PATCH 5/9] refactor(resourcetopo): extract timeout handling to separate methods Extract duplicate timeout logic into two helper methods: - executeHandlers(): for NodeHandler with timeout support - executeRelationHandlers(): for RelationHandler with timeout support This reduces code duplication and makes the main event handling logic cleaner. Co-Authored-By: Claude Opus 4.6 --- resourcetopo/event_queue.go | 188 +++++++++++++++--------------------- 1 file changed, 78 insertions(+), 110 deletions(-) diff --git a/resourcetopo/event_queue.go b/resourcetopo/event_queue.go index b06fbb4..86d8e61 100644 --- a/resourcetopo/event_queue.go +++ b/resourcetopo/event_queue.go @@ -76,83 +76,32 @@ func (m *manager) handleNodeEvent() { continue } - // Create context with timeout for all handlers in this event ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) - timeout := false - switch evtType { case EventTypeAdd: - for _, h := range storage.nodeUpdateHandler { - if timeout { - break - } - done := make(chan struct{}) - go func(handler NodeHandler) { - handler.OnAdd(node) - close(done) - }(h) - select { - case <-done: - // Handler completed - case <-ctx.Done(): - timeout = true - klog.Warningf("Node handler OnAdd timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - } - } + m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnAdd(node) + }, func() { + klog.Warningf("Node handler OnAdd timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }) case EventTypeUpdate: - for _, h := range storage.nodeUpdateHandler { - if timeout { - break - } - done := make(chan struct{}) - go func(handler NodeHandler) { - handler.OnUpdate(node) - close(done) - }(h) - select { - case <-done: - // Handler completed - case <-ctx.Done(): - timeout = true - klog.Warningf("Node handler OnUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - } - } + m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnUpdate(node) + }, func() { + klog.Warningf("Node handler OnUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }) case EventTypeDelete: - for _, h := range storage.nodeUpdateHandler { - if timeout { - break - } - done := make(chan struct{}) - go func(handler NodeHandler) { - handler.OnDelete(node) - close(done) - }(h) - select { - case <-done: - // Handler completed - case <-ctx.Done(): - timeout = true - klog.Warningf("Node handler OnDelete timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - } - } + m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnDelete(node) + }, func() { + klog.Warningf("Node handler OnDelete timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }) case EventTypeRelatedUpdate: - for _, h := range storage.nodeUpdateHandler { - if timeout { - break - } - done := make(chan struct{}) - go func(handler NodeHandler) { - handler.OnRelatedUpdate(node) - close(done) - }(h) - select { - case <-done: - // Handler completed - case <-ctx.Done(): - timeout = true - klog.Warningf("Node handler OnRelatedUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - } - } + m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnRelatedUpdate(node) + }, func() { + klog.Warningf("Node handler OnRelatedUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }) } cancel() m.nodeEventQueue.Done(item) @@ -190,55 +139,74 @@ func (m *manager) handleRelationEvent() { continue } - // Create context with timeout for all handlers in this event ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) - timeout := false - switch evtType { case EventTypeAdd: - for _, h := range handlers { - if timeout { - break - } - done := make(chan struct{}) - go func(handler RelationHandler) { - handler.OnAdd(preNode, postNode) - close(done) - }(h) - select { - case <-done: - // Handler completed - case <-ctx.Done(): - timeout = true - klog.Warningf("Relation handler OnAdd timeout for relation %s/%s -> %s/%s, skipping remaining handlers", - preNode.namespace, preNode.name, postNode.namespace, postNode.name) - } - } + m.executeRelationHandlers(ctx, handlers, func(h RelationHandler) { + h.OnAdd(preNode, postNode) + }, func() { + klog.Warningf("Relation handler OnAdd timeout for relation %s/%s -> %s/%s, skipping remaining handlers", + preNode.namespace, preNode.name, postNode.namespace, postNode.name) + }) case EventTypeDelete: - for _, h := range handlers { - if timeout { - break - } - done := make(chan struct{}) - go func(handler RelationHandler) { - handler.OnDelete(preNode, postNode) - close(done) - }(h) - select { - case <-done: - // Handler completed - case <-ctx.Done(): - timeout = true - klog.Warningf("Relation handler OnDelete timeout for relation %s/%s -> %s/%s, skipping remaining handlers", - preNode.namespace, preNode.name, postNode.namespace, postNode.name) - } - } + m.executeRelationHandlers(ctx, handlers, func(h RelationHandler) { + h.OnDelete(preNode, postNode) + }, func() { + klog.Warningf("Relation handler OnDelete timeout for relation %s/%s -> %s/%s, skipping remaining handlers", + preNode.namespace, preNode.name, postNode.namespace, postNode.name) + }) } cancel() m.relationEventQueue.Done(item) } } +// executeHandlers executes node handlers with timeout support. +// If timeout occurs, remaining handlers are skipped and onTimeout is called. +func (m *manager) executeHandlers(ctx context.Context, handlers []NodeHandler, exec func(h NodeHandler), onTimeout func()) { + timeout := false + for _, h := range handlers { + if timeout { + break + } + done := make(chan struct{}) + go func(handler NodeHandler) { + exec(handler) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + onTimeout() + } + } +} + +// executeRelationHandlers executes relation handlers with timeout support. +// If timeout occurs, remaining handlers are skipped and onTimeout is called. +func (m *manager) executeRelationHandlers(ctx context.Context, handlers []RelationHandler, exec func(h RelationHandler), onTimeout func()) { + timeout := false + for _, h := range handlers { + if timeout { + break + } + done := make(chan struct{}) + go func(handler RelationHandler) { + exec(handler) + close(done) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + onTimeout() + } + } +} + func (m *manager) newNodeEvent(info *nodeInfo, eType eventType) { key := m.encodeNodeEvent2String(eType, info) klog.V(6).Infof("New node event queue key: %v", key) From 76320d39568ae5406fc78085992e01e03156714d Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Sat, 25 Jul 2026 16:59:13 +0800 Subject: [PATCH 6/9] fix(resourcetopo): add context to handler interfaces, recover() guard, and timeout metrics - Add context.Context to NodeHandler and RelationHandler for cooperative cancellation. When EventHandlerTimeout expires, ctx is cancelled so handlers can select on ctx.Done() and return cleanly (no goroutine leak). - Add defer/recover() to executeHandlers and executeRelationHandlers to catch handler panics and prevent process crashes. Panics are logged and the done channel is always closed. - Add HandlerTimeoutMetrics() returning atomic counters for timed-out node and relation handler callbacks. - Fix checkGC() to use metaLock.RLock() instead of Lock() for a read-only boolean check, matching isObjectExisted(). - Revert GetClusterNode nil-return behavior change to preserve the original API contract for deleted-but-referenced nodes. - Update lock hierarchy docs: swap Level 4 (relationsLock) and Level 5 (lock) to match actual code ordering. - Update doc.go callback threading docs to reflect goroutine + timeout execution model. - Add tests: cooperative handler no-leak, panic recovery, goroutine leak detection, DirectRef delta updates, lock ordering stress. Co-Authored-By: Claude --- resourcetopo/doc.go | 23 +- resourcetopo/event_queue.go | 32 +- resourcetopo/examples/base_example.go | 12 +- resourcetopo/manager.go | 10 + resourcetopo/node_info.go | 4 +- resourcetopo/node_storage.go | 4 +- resourcetopo/resourcetopo_test.go | 508 ++++++++++++++++++++++++ resourcetopo/resourcetopo_utils_test.go | 50 ++- resourcetopo/types.go | 31 +- 9 files changed, 630 insertions(+), 44 deletions(-) diff --git a/resourcetopo/doc.go b/resourcetopo/doc.go index cb3895c..fb10e30 100644 --- a/resourcetopo/doc.go +++ b/resourcetopo/doc.go @@ -34,14 +34,14 @@ // - Protects: clusterNodes map // - Note: Never acquire while holding any nodeInfo.lock // -// Level 4: nodeInfo.lock +// Level 4: nodeInfo.relationsLock +// - Protects: labelRelations +// - Note: On the same node, must be acquired before nodeInfo.lock +// +// Level 5: nodeInfo.lock // - Protects: PreOrders/PostOrders relation lists (directReferredPreOrders, etc.) // - Rule: Always acquire preNode.lock before postNode.lock // -// Level 5: nodeInfo.relationsLock -// - Protects: labelRelations -// - Rule: Acquire before metaLock if both needed -// // Level 6: nodeInfo.metaLock // - Protects: labels, ownerNodes, objectExisted // - This is the innermost lock - never acquire other locks while holding metaLock @@ -80,7 +80,14 @@ // 4. Callback Threads (External - via Event Processors): // - Source: User-provided NodeHandler and RelationHandler callbacks // - Entry: Called by handleNodeEvent/handleRelationEvent -// - Lock behavior: No locks held during callback execution +// - Execution: Each callback runs in its own goroutine with a configurable timeout +// (ManagerConfig.EventHandlerTimeout, default 1s). If a callback does not return +// within the timeout, a warning is logged and all remaining handlers for that +// event are skipped. +// - Lock behavior: No resourcetopo locks held during callback execution +// - Concurrency: Multiple callbacks (across different handlers or events) may run +// concurrently. Handler implementations that share mutable state must be +// thread-safe. // - Safety: Handlers receive node references but should not cache them long-term // // Thread Safety Guidelines @@ -88,7 +95,9 @@ // - NodeInfo references returned by GetNode() are safe for concurrent reads // - Do NOT call AddTopologyConfig() after Start() - not thread safe // - Handler registration (AddNodeHandler, AddRelationHandler) is thread safe -// - All callbacks (OnAdd, OnUpdate, etc.) are called synchronously - do not block +// - Callbacks (OnAdd, OnUpdate, etc.) run in goroutines with a timeout. +// Do not block indefinitely — timed-out callbacks cause remaining handlers +// to be skipped. Implementations must be thread-safe if they share state. // - nodeInfo.lock can be held for extended periods during relation changes // - metaLock is never held during callbacks or cross-node operations // diff --git a/resourcetopo/event_queue.go b/resourcetopo/event_queue.go index 86d8e61..7f127b2 100644 --- a/resourcetopo/event_queue.go +++ b/resourcetopo/event_queue.go @@ -80,26 +80,30 @@ func (m *manager) handleNodeEvent() { switch evtType { case EventTypeAdd: m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { - h.OnAdd(node) + h.OnAdd(ctx, node) }, func() { + m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnAdd timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) case EventTypeUpdate: m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { - h.OnUpdate(node) + h.OnUpdate(ctx, node) }, func() { + m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) case EventTypeDelete: m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { - h.OnDelete(node) + h.OnDelete(ctx, node) }, func() { + m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnDelete timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) case EventTypeRelatedUpdate: m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { - h.OnRelatedUpdate(node) + h.OnRelatedUpdate(ctx, node) }, func() { + m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnRelatedUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) } @@ -143,15 +147,17 @@ func (m *manager) handleRelationEvent() { switch evtType { case EventTypeAdd: m.executeRelationHandlers(ctx, handlers, func(h RelationHandler) { - h.OnAdd(preNode, postNode) + h.OnAdd(ctx, preNode, postNode) }, func() { + m.relationHandlerTimeouts.Add(1) klog.Warningf("Relation handler OnAdd timeout for relation %s/%s -> %s/%s, skipping remaining handlers", preNode.namespace, preNode.name, postNode.namespace, postNode.name) }) case EventTypeDelete: m.executeRelationHandlers(ctx, handlers, func(h RelationHandler) { - h.OnDelete(preNode, postNode) + h.OnDelete(ctx, preNode, postNode) }, func() { + m.relationHandlerTimeouts.Add(1) klog.Warningf("Relation handler OnDelete timeout for relation %s/%s -> %s/%s, skipping remaining handlers", preNode.namespace, preNode.name, postNode.namespace, postNode.name) }) @@ -171,8 +177,13 @@ func (m *manager) executeHandlers(ctx context.Context, handlers []NodeHandler, e } done := make(chan struct{}) go func(handler NodeHandler) { + defer func() { + if r := recover(); r != nil { + klog.Errorf("Node handler panicked: %v, treating as failed", r) + } + close(done) + }() exec(handler) - close(done) }(h) select { case <-done: @@ -194,8 +205,13 @@ func (m *manager) executeRelationHandlers(ctx context.Context, handlers []Relati } done := make(chan struct{}) go func(handler RelationHandler) { + defer func() { + if r := recover(); r != nil { + klog.Errorf("Relation handler panicked: %v, treating as failed", r) + } + close(done) + }() exec(handler) - close(done) }(h) select { case <-done: diff --git a/resourcetopo/examples/base_example.go b/resourcetopo/examples/base_example.go index ca093ed..4bfa6da 100644 --- a/resourcetopo/examples/base_example.go +++ b/resourcetopo/examples/base_example.go @@ -224,19 +224,19 @@ var _ resourcetopo.NodeHandler = &podEventhandler{} type podEventhandler struct{} -func (p *podEventhandler) OnAdd(info resourcetopo.NodeInfo) { +func (p *podEventhandler) OnAdd(ctx context.Context, info resourcetopo.NodeInfo) { klog.Infof("received add event for pod %s", info.NodeInfo().String()) } -func (p *podEventhandler) OnUpdate(info resourcetopo.NodeInfo) { +func (p *podEventhandler) OnUpdate(ctx context.Context, info resourcetopo.NodeInfo) { klog.Infof("received update event for pod %s", info.NodeInfo().String()) } -func (p *podEventhandler) OnDelete(info resourcetopo.NodeInfo) { +func (p *podEventhandler) OnDelete(ctx context.Context, info resourcetopo.NodeInfo) { klog.Infof("received delete event for pod %s", info.NodeInfo().String()) } -func (p *podEventhandler) OnRelatedUpdate(info resourcetopo.NodeInfo) { +func (p *podEventhandler) OnRelatedUpdate(ctx context.Context, info resourcetopo.NodeInfo) { klog.Infof("related node has changed and effected pod %s", info.NodeInfo().String()) klog.Infof("related pre nodes are %s", nodes2Str(info.GetPreOrders())) klog.Infof("related post nodes are %s", nodes2Str(info.GetPostOrders())) @@ -246,11 +246,11 @@ var _ resourcetopo.RelationHandler = &appPodRelationEventHandler{} type appPodRelationEventHandler struct{} -func (s *appPodRelationEventHandler) OnAdd(preOrder, postOrder resourcetopo.NodeInfo) { +func (s *appPodRelationEventHandler) OnAdd(ctx context.Context, preOrder, postOrder resourcetopo.NodeInfo) { klog.Infof("received relation add event for %s -> %s", preOrder.NodeInfo().String(), postOrder.NodeInfo().String()) } -func (s *appPodRelationEventHandler) OnDelete(preOrder, postOrder resourcetopo.NodeInfo) { +func (s *appPodRelationEventHandler) OnDelete(ctx context.Context, preOrder, postOrder resourcetopo.NodeInfo) { klog.Infof("received relation delete event for %s -> %s", preOrder.NodeInfo().String(), postOrder.NodeInfo().String()) } diff --git a/resourcetopo/manager.go b/resourcetopo/manager.go index 578e6c5..d66015b 100644 --- a/resourcetopo/manager.go +++ b/resourcetopo/manager.go @@ -20,6 +20,7 @@ import ( "container/list" "fmt" "sync" + "sync/atomic" "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -40,6 +41,15 @@ type manager struct { storagesLock sync.RWMutex // protects storages map eventHandlerTimeout time.Duration // timeout for event handler callbacks + + // metrics + nodeHandlerTimeouts atomic.Int64 + relationHandlerTimeouts atomic.Int64 +} + +// HandlerTimeoutMetrics returns the count of timed-out node and relation handlers. +func (m *manager) HandlerTimeoutMetrics() (nodeTimeouts, relationTimeouts int64) { + return m.nodeHandlerTimeouts.Load(), m.relationHandlerTimeouts.Load() } func NewResourcesTopoManager(cfg ManagerConfig) (Manager, error) { diff --git a/resourcetopo/node_info.go b/resourcetopo/node_info.go index 23364d9..50ce7a9 100644 --- a/resourcetopo/node_info.go +++ b/resourcetopo/node_info.go @@ -258,9 +258,9 @@ func (n *nodeInfo) checkGC() { } // Then acquire metaLock to check object existence (Level 6) - n.metaLock.Lock() + n.metaLock.RLock() objectExisted := n.objectExisted - n.metaLock.Unlock() + n.metaLock.RUnlock() // n maybe set as notExisted, but still have related nodes to handle if !objectExisted { diff --git a/resourcetopo/node_storage.go b/resourcetopo/node_storage.go index 3823e61..dad2c52 100644 --- a/resourcetopo/node_storage.go +++ b/resourcetopo/node_storage.go @@ -90,9 +90,7 @@ func (s *nodeStorage) GetNode(namespacedName types.NamespacedName) (NodeInfo, er func (s *nodeStorage) GetClusterNode(cluster string, namespacedName types.NamespacedName) (NodeInfo, error) { node := s.getNode(cluster, namespacedName.Namespace, namespacedName.Name) - if node == nil || !node.isObjectExisted() { - // return nil rather than node ref here, since the node may be a virtual node - // todo add a unittest + if node != nil && !node.isObjectExisted() { return nil, nil } return node, nil diff --git a/resourcetopo/resourcetopo_test.go b/resourcetopo/resourcetopo_test.go index 85f2d92..35d872e 100644 --- a/resourcetopo/resourcetopo_test.go +++ b/resourcetopo/resourcetopo_test.go @@ -30,6 +30,7 @@ import ( "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes/fake" + "runtime" "k8s.io/klog/v2" ) @@ -2022,5 +2023,512 @@ var _ = Describe("test suite for multi routine", func() { g.Expect(len(rsNode.GetPostOrders())).To(Equal(totalNum / len(rsNames))) } }, "10s").Should(Succeed()) + }) + }) + +var _ = Describe("test handler timeout protection", func() { + var manager Manager + var fakeClient *fake.Clientset + var ctx context.Context + var cancel func() + + BeforeEach(func() { + fakeClient = fake.NewSimpleClientset() + k8sInformerFactory := informers.NewSharedInformerFactory(fakeClient, 0) + var err error + cfg := buildManagerConfig(buildClusterTest(k8sInformerFactory)) + cfg.EventHandlerTimeout = 100 * time.Millisecond + manager, err = NewResourcesTopoManager(*cfg) + Expect(err).NotTo(HaveOccurred()) + + ctx, cancel = context.WithCancel(context.Background()) + manager.Start(ctx.Done()) + k8sInformerFactory.Start(ctx.Done()) + k8sInformerFactory.WaitForCacheSync(ctx.Done()) + }) + AfterEach(func() { + cancel() + }) + + It("blocking node handler should timeout and not block the event loop", func() { + ns := "testnodetimeout" + crName := "crName" + saName := "saName" + + blockCh := make(chan struct{}) + blockingHandler := &blockingNodeHandler{blockCh: blockCh} + Expect(manager.AddNodeHandler(ClusterRoleMeta, blockingHandler)).To(Succeed()) + + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + time.Sleep(300 * time.Millisecond) + + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName), metav1.CreateOptions{})).NotTo(BeNil()) + + Eventually(func(g Gomega) { + g.Expect(saHandler.matchExpected()).To(BeTrue()) + }).Should(Succeed()) + }) + + It("blocking relation handler should timeout and not block the event loop", func() { + ns := "testreltimeout" + crName := "crName" + crbName := "crbtest" + saName := "saName" + + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + crHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleMeta, crHandler)).To(Succeed()) + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName), metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(func() bool { return saHandler.matchExpected() }) + + crHandler.addCallExpected() + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(func() bool { return crHandler.matchExpected() }) + + blockCh := make(chan struct{}) + blockingRelHandler := &blockingRelationHandler{blockCh: blockCh} + Expect(manager.AddRelationHandler(ClusterRoleBindingMeta, ServiceAccountMeta, blockingRelHandler)).To(Succeed()) + + crbHandler.addCallExpected() + Expect(fakeClient.RbacV1().ClusterRoleBindings().Create(ctx, newClusterRoleBinding(crbName, crName, + []types.NamespacedName{{Name: saName, Namespace: ns}}), + metav1.CreateOptions{})).NotTo(BeNil()) + time.Sleep(300 * time.Millisecond) + + saName2 := "saName2" + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName2), metav1.CreateOptions{})).NotTo(BeNil()) + + Eventually(func(g Gomega) { + g.Expect(saHandler.matchExpected()).To(BeTrue()) + }).Should(Succeed()) + }) +}) + +// blockingNodeHandler blocks forever on every callback, used for timeout testing +type blockingNodeHandler struct { + blockCh chan struct{} +} + +func (h *blockingNodeHandler) OnAdd(ctx context.Context, info NodeInfo) { + select { + case <-h.blockCh: + case <-ctx.Done(): + } +} +func (h *blockingNodeHandler) OnUpdate(ctx context.Context, info NodeInfo) { + select { + case <-h.blockCh: + case <-ctx.Done(): + } +} +func (h *blockingNodeHandler) OnDelete(ctx context.Context, info NodeInfo) { + select { + case <-h.blockCh: + case <-ctx.Done(): + } +} +func (h *blockingNodeHandler) OnRelatedUpdate(ctx context.Context, info NodeInfo) { + select { + case <-h.blockCh: + case <-ctx.Done(): + } +} + +var _ NodeHandler = &blockingNodeHandler{} + +// blockingRelationHandler blocks forever on every relation callback, used for timeout testing +type blockingRelationHandler struct { + blockCh chan struct{} +} + +func (h *blockingRelationHandler) OnAdd(ctx context.Context, preNode, postNode NodeInfo) { + select { + case <-h.blockCh: + case <-ctx.Done(): + } +} +func (h *blockingRelationHandler) OnDelete(ctx context.Context, preNode, postNode NodeInfo) { + select { + case <-h.blockCh: + case <-ctx.Done(): + } +} + +var _ RelationHandler = &blockingRelationHandler{} + +var _ = Describe("test DirectRef delta-based relation updates", func() { + var manager Manager + var fakeClient *fake.Clientset + var clusterRoleBindingHandler, clusterRoleHandler, saHandler *objecthandler + var saBindingRelation, roleBindingRelation *relationHandler + + var clusterRoleBindingStorage, clusterroleStorage, saStorage TopoNodeStorage + var ctx context.Context + var cancel func() + + checkAll := func() bool { + return clusterRoleBindingHandler.matchExpected() && + clusterRoleHandler.matchExpected() && + saHandler.matchExpected() && + saBindingRelation.matchExpected() && + roleBindingRelation.matchExpected() + } + + BeforeEach(func() { + fakeClient = fake.NewSimpleClientset() + k8sInformerFactory := informers.NewSharedInformerFactory(fakeClient, 0) + var err error + manager, err = NewResourcesTopoManager(*buildManagerConfig(buildClusterTest(k8sInformerFactory))) + Expect(err).NotTo(HaveOccurred()) + + clusterRoleBindingStorage, _ = manager.GetTopoNodeStorage(ClusterRoleBindingMeta) + clusterroleStorage, _ = manager.GetTopoNodeStorage(ClusterRoleMeta) + saStorage, _ = manager.GetTopoNodeStorage(ServiceAccountMeta) + + Expect(clusterRoleBindingStorage).NotTo(BeNil()) + Expect(clusterroleStorage).NotTo(BeNil()) + Expect(saStorage).NotTo(BeNil()) + + ctx, cancel = context.WithCancel(context.Background()) + clusterRoleBindingHandler = &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, clusterRoleBindingHandler)).NotTo(HaveOccurred()) + clusterRoleHandler = &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleMeta, clusterRoleHandler)).NotTo(HaveOccurred()) + saHandler = &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).NotTo(HaveOccurred()) + + roleBindingRelation = &relationHandler{} + Expect(manager.AddRelationHandler(ClusterRoleBindingMeta, ClusterRoleMeta, roleBindingRelation)).To(Succeed()) + saBindingRelation = &relationHandler{} + Expect(manager.AddRelationHandler(ClusterRoleBindingMeta, ServiceAccountMeta, saBindingRelation)).To(Succeed()) + + manager.Start(ctx.Done()) + k8sInformerFactory.Start(ctx.Done()) + k8sInformerFactory.WaitForCacheSync(ctx.Done()) + }) + AfterEach(func() { + cancel() + }) + + It("updating DirectRefs should only fire callbacks for changed refs", func() { + ns := "testdeltadir" + crbName := "crbtest" + crName := "crName" + saName := "saName" + saName2 := "saName2" + saName3 := "saName3" + + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName), metav1.CreateOptions{})).NotTo(BeNil()) + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName2), metav1.CreateOptions{})).NotTo(BeNil()) + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName3), metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + clusterRoleHandler.addCallExpected() + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + clusterRoleBindingHandler.addCallExpected() + saBindingRelation.addCallExpected() // saName + saBindingRelation.addCallExpected() // saName2 + roleBindingRelation.addCallExpected() + Expect(fakeClient.RbacV1().ClusterRoleBindings().Create(ctx, newClusterRoleBinding(crbName, crName, + []types.NamespacedName{{Name: saName, Namespace: ns}, {Name: saName2, Namespace: ns}}), + metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + // Update CRB: replace saName2 with saName3 + // Delta: only saName3 added + saName2 deleted (not 2+2) + saBindingRelation.addCallExpected() // saName3 added + saBindingRelation.deleteCallExpected() // saName2 deleted + clusterRoleBindingHandler.updateCallExpected() + Expect(fakeClient.RbacV1().ClusterRoleBindings().Update(ctx, newClusterRoleBinding(crbName, crName, + []types.NamespacedName{{Name: saName, Namespace: ns}, {Name: saName3, Namespace: ns}}), + metav1.UpdateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + Eventually(func(g Gomega) { + crb, _ := clusterRoleBindingStorage.GetNode(types.NamespacedName{Name: crbName}) + g.Expect(crb).NotTo(BeNil()) + g.Expect(crb.GetPostOrders()).To(HaveLen(3)) + g.Expect(checkAll()).To(BeTrue()) + }).Should(Succeed()) + }) + + It("adding multiple and removing all DirectRefs should work", func() { + ns := "testdeltaall" + crbName := "crbtest" + crName := "crName" + saName := "saName" + saName2 := "saName2" + + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName), metav1.CreateOptions{})).NotTo(BeNil()) + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName2), metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + clusterRoleHandler.addCallExpected() + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + clusterRoleBindingHandler.addCallExpected() + saBindingRelation.addCallExpected() + saBindingRelation.addCallExpected() + roleBindingRelation.addCallExpected() + Expect(fakeClient.RbacV1().ClusterRoleBindings().Create(ctx, newClusterRoleBinding(crbName, crName, + []types.NamespacedName{{Name: saName, Namespace: ns}, {Name: saName2, Namespace: ns}}), + metav1.CreateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + saBindingRelation.deleteCallExpected() + saBindingRelation.deleteCallExpected() + clusterRoleBindingHandler.updateCallExpected() + Expect(fakeClient.RbacV1().ClusterRoleBindings().Update(ctx, newClusterRoleBinding(crbName, crName, + []types.NamespacedName{}), + metav1.UpdateOptions{})).NotTo(BeNil()) + syncStatus(checkAll) + + Eventually(func(g Gomega) { + crb, _ := clusterRoleBindingStorage.GetNode(types.NamespacedName{Name: crbName}) + g.Expect(crb).NotTo(BeNil()) + g.Expect(crb.GetPostOrders()).To(HaveLen(1)) + g.Expect(checkAll()).To(BeTrue()) + }).Should(Succeed()) + }) +}) + +var _ = Describe("test concurrent DirectRef operations for lock ordering", func() { + var manager Manager + var fakeClient *fake.Clientset + var ctx context.Context + var cancel func() + + BeforeEach(func() { + fakeClient = fake.NewSimpleClientset() + k8sInformerFactory := informers.NewSharedInformerFactory(fakeClient, 0) + var err error + manager, err = NewResourcesTopoManager(*buildManagerConfig(buildClusterTest(k8sInformerFactory))) + Expect(err).NotTo(HaveOccurred()) + + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + crHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleMeta, crHandler)).To(Succeed()) + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + Expect(manager.AddRelationHandler(ClusterRoleBindingMeta, ClusterRoleMeta, &relationHandler{})).To(Succeed()) + Expect(manager.AddRelationHandler(ClusterRoleBindingMeta, ServiceAccountMeta, &relationHandler{})).To(Succeed()) + + ctx, cancel = context.WithCancel(context.Background()) + manager.Start(ctx.Done()) + k8sInformerFactory.Start(ctx.Done()) + k8sInformerFactory.WaitForCacheSync(ctx.Done()) + }) + AfterEach(func() { + cancel() + }) + + It("concurrent create and delete should not deadlock", func() { + ns := "testlockorder" + crName := "crName" + numResources := 20 + + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + + var wg sync.WaitGroup + + for i := 0; i < numResources; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + saName := fmt.Sprintf("sa-%d", idx) + fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName), metav1.CreateOptions{}) + }(i) + } + wg.Wait() + time.Sleep(500 * time.Millisecond) + + for i := 0; i < numResources; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + crbName := fmt.Sprintf("crb-%d", idx) + saName := fmt.Sprintf("sa-%d", idx) + fakeClient.RbacV1().ClusterRoleBindings().Create(ctx, newClusterRoleBinding(crbName, crName, + []types.NamespacedName{{Name: saName, Namespace: ns}}), + metav1.CreateOptions{}) + }(i) + + if i%3 == 0 { + wg.Add(1) + go func(idx int) { + defer wg.Done() + saName := fmt.Sprintf("sa-%d", idx) + fakeClient.CoreV1().ServiceAccounts(ns).Delete(ctx, saName, metav1.DeleteOptions{}) + }(i) + } + } + wg.Wait() + + for i := 0; i < numResources; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + crbName := fmt.Sprintf("crb-%d", idx) + fakeClient.RbacV1().ClusterRoleBindings().Delete(ctx, crbName, metav1.DeleteOptions{}) + }(i) + if i%2 == 0 { + wg.Add(1) + go func(idx int) { + defer wg.Done() + saName := fmt.Sprintf("sa-%d", idx) + fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName), metav1.CreateOptions{}) + }(i) + } + } + wg.Wait() + + time.Sleep(500 * time.Millisecond) + crbStorage, _ := manager.GetTopoNodeStorage(ClusterRoleBindingMeta) + Expect(crbStorage).NotTo(BeNil()) + }) +}) + +// --- Tests for known issues: goroutine leak and panic absorption --- + +var _ = Describe("test goroutine leak and panic handling", func() { + var manager Manager + var fakeClient *fake.Clientset + var ctx context.Context + var cancel func() + + BeforeEach(func() { + fakeClient = fake.NewSimpleClientset() + k8sInformerFactory := informers.NewSharedInformerFactory(fakeClient, 0) + var err error + cfg := buildManagerConfig(buildClusterTest(k8sInformerFactory)) + cfg.EventHandlerTimeout = 100 * time.Millisecond + manager, err = NewResourcesTopoManager(*cfg) + Expect(err).NotTo(HaveOccurred()) + + ctx, cancel = context.WithCancel(context.Background()) + manager.Start(ctx.Done()) + k8sInformerFactory.Start(ctx.Done()) + k8sInformerFactory.WaitForCacheSync(ctx.Done()) + }) + AfterEach(func() { + cancel() + // Give goroutines time to clean up + time.Sleep(200 * time.Millisecond) + }) + + It("cooperative handler does not leak goroutine on timeout", func() { + crName := "crNoLeak" + + // Register a cooperative blocking handler for ClusterRole + // (selects on ctx.Done() so it returns when context is cancelled) + blockCh := make(chan struct{}) + blockingHandler := &blockingNodeHandler{blockCh: blockCh} + Expect(manager.AddNodeHandler(ClusterRoleMeta, blockingHandler)).To(Succeed()) + + // Register normal handlers for other types + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + + // Let the event loop stabilize + time.Sleep(100 * time.Millisecond) + + // Count baseline goroutines + baseline := countStableGoroutines(3, 100*time.Millisecond) + + // Trigger handler — ctx will be cancelled after timeout, handler returns + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + + // Wait well past timeout for handler to return and goroutine to exit + time.Sleep(500 * time.Millisecond) + + // Count goroutines after — should be stable (no leak) + afterLeak := countStableGoroutines(3, 100*time.Millisecond) + + // Cooperative handler exits on ctx.Done(), goroutine count should be stable + Expect(afterLeak).To(BeNumerically("<=", baseline+1), + fmt.Sprintf("Expected no goroutine leak: baseline=%d, after=%d", baseline, afterLeak)) + }) + + It("handler panic is silently absorbed as timeout", func() { + // When a handler panics, executeHandlers' recover() catches it, + // logs an error, and closes the done channel. + // This test verifies the panic does NOT crash the process + // and the event loop continues processing subsequent events. + + ns := "testpanic" + saName := "saAfterPanic" + + // Register a panicking handler for ClusterRole + panickingHandler := &panickingNodeHandler{} + Expect(manager.AddNodeHandler(ClusterRoleMeta, panickingHandler)).To(Succeed()) + + // Register normal handlers + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + + // Trigger the panicking handler — it panics in a goroutine, + // the panic is absorbed, and the event is treated as timed out + crName := "crPanic" + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + + // Wait for the timeout (100ms + buffer) + time.Sleep(300 * time.Millisecond) + + // Verify the event loop is still alive + saHandler.addCallExpected() + Expect(fakeClient.CoreV1().ServiceAccounts(ns).Create(ctx, newServiceAccount(ns, saName), metav1.CreateOptions{})).NotTo(BeNil()) + + Eventually(func(g Gomega) { + g.Expect(saHandler.matchExpected()).To(BeTrue()) + }).Should(Succeed()) }) }) + +// countStableGoroutines samples NumGoroutine multiple times and returns +// the minimum observed count (to filter out transient goroutines). +func countStableGoroutines(samples int, interval time.Duration) int { + min := runtime.NumGoroutine() + for i := 1; i < samples; i++ { + time.Sleep(interval) + n := runtime.NumGoroutine() + if n < min { + min = n + } + } + return min +} + +// panickingNodeHandler panics on every callback, used to test panic absorption. +type panickingNodeHandler struct{} + +func (h *panickingNodeHandler) OnAdd(ctx context.Context, info NodeInfo) { panic("intentional panic in OnAdd") } +func (h *panickingNodeHandler) OnUpdate(ctx context.Context, info NodeInfo) { panic("intentional panic in OnUpdate") } +func (h *panickingNodeHandler) OnDelete(ctx context.Context, info NodeInfo) { panic("intentional panic in OnDelete") } +func (h *panickingNodeHandler) OnRelatedUpdate(ctx context.Context, info NodeInfo) { panic("intentional panic in OnRelatedUpdate") } + +var _ NodeHandler = &panickingNodeHandler{} diff --git a/resourcetopo/resourcetopo_utils_test.go b/resourcetopo/resourcetopo_utils_test.go index d09b6ea..d3b411c 100644 --- a/resourcetopo/resourcetopo_utils_test.go +++ b/resourcetopo/resourcetopo_utils_test.go @@ -17,8 +17,10 @@ package resourcetopo import ( + "context" "encoding/json" "fmt" + "sync" . "github.com/onsi/ginkgo" @@ -509,6 +511,8 @@ func getMultiClusterDepend(object metav1.Object) []MultiClusterDepend { var _ NodeHandler = &objecthandler{} type objecthandler struct { + mu sync.Mutex + addCounter int updateCounter int relatedCounter int @@ -521,47 +525,65 @@ type objecthandler struct { // change loglevel flag to 0 to enable log output const loglevel = 1 -func (o *objecthandler) OnAdd(info NodeInfo) { +func (o *objecthandler) OnAdd(ctx context.Context, info NodeInfo) { klog.V(loglevel).Infof("received added object %v %v", info.TypeInfo(), info.NodeInfo()) + o.mu.Lock() o.addCounter-- + o.mu.Unlock() o.rangeNode(info) } -func (o *objecthandler) OnUpdate(info NodeInfo) { +func (o *objecthandler) OnUpdate(ctx context.Context, info NodeInfo) { klog.V(loglevel).Infof("received updated object %v %v", info.TypeInfo(), info.NodeInfo()) + o.mu.Lock() o.updateCounter-- + o.mu.Unlock() o.rangeNode(info) } -func (o *objecthandler) OnDelete(info NodeInfo) { +func (o *objecthandler) OnDelete(ctx context.Context, info NodeInfo) { klog.V(loglevel).Infof("received deleted object %v %v", info.TypeInfo(), info.NodeInfo()) + o.mu.Lock() o.deletedCounter-- + o.mu.Unlock() o.rangeNode(info) } -func (o *objecthandler) OnRelatedUpdate(info NodeInfo) { +func (o *objecthandler) OnRelatedUpdate(ctx context.Context, info NodeInfo) { klog.V(loglevel).Infof("received related updated object %v %v", info.TypeInfo(), info.NodeInfo()) + o.mu.Lock() o.relatedCounter-- + o.mu.Unlock() o.rangeNode(info) } func (o *objecthandler) addCallExpected() { + o.mu.Lock() o.addCounter++ + o.mu.Unlock() } func (o *objecthandler) updateCallExpected() { + o.mu.Lock() o.updateCounter++ + o.mu.Unlock() } func (o *objecthandler) deleteCallExpected() { + o.mu.Lock() o.deletedCounter++ + o.mu.Unlock() } func (o *objecthandler) relatedCallExpected() { + o.mu.Lock() o.relatedCounter++ + o.mu.Unlock() } func (h *objecthandler) matchExpected() bool { + h.mu.Lock() + defer h.mu.Unlock() return h.addCounter == 0 && h.updateCounter == 0 && h.deletedCounter == 0 && h.relatedCounter == 0 } @@ -591,40 +613,56 @@ func (o *objecthandler) rangePostOrder(node NodeInfo) { } func (o *objecthandler) string() string { + o.mu.Lock() + defer o.mu.Unlock() return fmt.Sprintf("{add: %d, update: %d, delete %d, relatedUpdate %d}", o.addCounter, o.updateCounter, o.deletedCounter, o.relatedCounter) } var _ RelationHandler = &relationHandler{} type relationHandler struct { + mu sync.Mutex + addCounter int deleteCounter int } -func (r *relationHandler) OnAdd(preNode, postNode NodeInfo) { +func (r *relationHandler) OnAdd(ctx context.Context, preNode, postNode NodeInfo) { klog.V(loglevel).Infof("received added relation, preNode %v %v, postNode %v %v", preNode.TypeInfo(), preNode.NodeInfo(), postNode.TypeInfo(), postNode.NodeInfo()) + r.mu.Lock() r.addCounter-- + r.mu.Unlock() } -func (r *relationHandler) OnDelete(preNode, postNode NodeInfo) { +func (r *relationHandler) OnDelete(ctx context.Context, preNode, postNode NodeInfo) { klog.V(loglevel).Infof("received deleted relation, preNode %v %v, postNode %v %v", preNode.TypeInfo(), preNode.NodeInfo(), postNode.TypeInfo(), postNode.NodeInfo()) + r.mu.Lock() r.deleteCounter-- + r.mu.Unlock() } func (r *relationHandler) addCallExpected() { + r.mu.Lock() r.addCounter++ + r.mu.Unlock() } func (r *relationHandler) deleteCallExpected() { + r.mu.Lock() r.deleteCounter++ + r.mu.Unlock() } func (r *relationHandler) matchExpected() bool { + r.mu.Lock() + defer r.mu.Unlock() return r.addCounter == 0 && r.deleteCounter == 0 } func (r *relationHandler) string() string { + r.mu.Lock() + defer r.mu.Unlock() return fmt.Sprintf("{add: %d, delete %d}", r.addCounter, r.deleteCounter) } diff --git a/resourcetopo/types.go b/resourcetopo/types.go index 91f42ef..db2b7a8 100644 --- a/resourcetopo/types.go +++ b/resourcetopo/types.go @@ -17,6 +17,7 @@ package resourcetopo import ( + "context" "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -170,26 +171,32 @@ type ResourceRelation struct { } type RelationHandler interface { - // OnAdd is called after relation newly added between preOrder and postOrder - OnAdd(preOrder, postOrder NodeInfo) + // OnAdd is called after relation newly added between preOrder and postOrder. + // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + OnAdd(ctx context.Context, preOrder, postOrder NodeInfo) - // OnDelete is called after relation deleted between preOrder and postOrder - OnDelete(preOrder, postOrder NodeInfo) + // OnDelete is called after relation deleted between preOrder and postOrder. + // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + OnDelete(ctx context.Context, preOrder, postOrder NodeInfo) } type NodeHandler interface { - // OnAdd is called when the node related resource object is added - OnAdd(info NodeInfo) + // OnAdd is called when the node related resource object is added. + // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + OnAdd(ctx context.Context, info NodeInfo) - // OnUpdate is called when the node related resource object is updated - OnUpdate(info NodeInfo) + // OnUpdate is called when the node related resource object is updated. + // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + OnUpdate(ctx context.Context, info NodeInfo) - // OnDelete is called when the node related resource object is deleted - OnDelete(info NodeInfo) + // OnDelete is called when the node related resource object is deleted. + // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + OnDelete(ctx context.Context, info NodeInfo) // OnRelatedUpdate is called when the nodes' - // post-order(or pre-order if ReverseNotice configured) object is updated - OnRelatedUpdate(info NodeInfo) + // post-order(or pre-order if ReverseNotice configured) object is updated. + // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + OnRelatedUpdate(ctx context.Context, info NodeInfo) } // Object define get namespace/name/labels for now, From 078044d95939656de5722b02067bfd2500f045fd Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Sat, 25 Jul 2026 19:27:22 +0800 Subject: [PATCH 7/9] refactor(resourcetopo): unify executeHandlers with generics, use struct{} for map set - Replace executeHandlers/executeRelationHandlers (40-line copy-paste) with a single generic executeHandlers[H any]() function. - Change map[nodeMeta]interface{} to map[nodeMeta]struct{} for idiomatic Go set semantics (zero-value sentinel, zero allocation). Co-Authored-By: Claude --- resourcetopo/event_queue.go | 63 ++++++++------------------- resourcetopo/node_topology_updater.go | 6 +-- 2 files changed, 21 insertions(+), 48 deletions(-) diff --git a/resourcetopo/event_queue.go b/resourcetopo/event_queue.go index 7f127b2..aa1d797 100644 --- a/resourcetopo/event_queue.go +++ b/resourcetopo/event_queue.go @@ -79,30 +79,30 @@ func (m *manager) handleNodeEvent() { ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) switch evtType { case EventTypeAdd: - m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnAdd(ctx, node) - }, func() { + }, "Node handler", func() { m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnAdd timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) case EventTypeUpdate: - m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnUpdate(ctx, node) - }, func() { + }, "Node handler", func() { m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) case EventTypeDelete: - m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnDelete(ctx, node) - }, func() { + }, "Node handler", func() { m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnDelete timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) case EventTypeRelatedUpdate: - m.executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnRelatedUpdate(ctx, node) - }, func() { + }, "Node handler", func() { m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnRelatedUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) }) @@ -146,17 +146,17 @@ func (m *manager) handleRelationEvent() { ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) switch evtType { case EventTypeAdd: - m.executeRelationHandlers(ctx, handlers, func(h RelationHandler) { + executeHandlers(ctx, handlers, func(h RelationHandler) { h.OnAdd(ctx, preNode, postNode) - }, func() { + }, "Relation handler", func() { m.relationHandlerTimeouts.Add(1) klog.Warningf("Relation handler OnAdd timeout for relation %s/%s -> %s/%s, skipping remaining handlers", preNode.namespace, preNode.name, postNode.namespace, postNode.name) }) case EventTypeDelete: - m.executeRelationHandlers(ctx, handlers, func(h RelationHandler) { + executeHandlers(ctx, handlers, func(h RelationHandler) { h.OnDelete(ctx, preNode, postNode) - }, func() { + }, "Relation handler", func() { m.relationHandlerTimeouts.Add(1) klog.Warningf("Relation handler OnDelete timeout for relation %s/%s -> %s/%s, skipping remaining handlers", preNode.namespace, preNode.name, postNode.namespace, postNode.name) @@ -167,47 +167,20 @@ func (m *manager) handleRelationEvent() { } } -// executeHandlers executes node handlers with timeout support. -// If timeout occurs, remaining handlers are skipped and onTimeout is called. -func (m *manager) executeHandlers(ctx context.Context, handlers []NodeHandler, exec func(h NodeHandler), onTimeout func()) { +// executeHandlers executes handlers with timeout support using generics. +// Each handler runs in its own goroutine. If the context deadline fires before +// a handler completes, remaining handlers are skipped and onTimeout is called. +func executeHandlers[H any](ctx context.Context, handlers []H, exec func(H), panicMsg string, onTimeout func()) { timeout := false for _, h := range handlers { if timeout { break } done := make(chan struct{}) - go func(handler NodeHandler) { + go func(handler H) { defer func() { if r := recover(); r != nil { - klog.Errorf("Node handler panicked: %v, treating as failed", r) - } - close(done) - }() - exec(handler) - }(h) - select { - case <-done: - // Handler completed - case <-ctx.Done(): - timeout = true - onTimeout() - } - } -} - -// executeRelationHandlers executes relation handlers with timeout support. -// If timeout occurs, remaining handlers are skipped and onTimeout is called. -func (m *manager) executeRelationHandlers(ctx context.Context, handlers []RelationHandler, exec func(h RelationHandler), onTimeout func()) { - timeout := false - for _, h := range handlers { - if timeout { - break - } - done := make(chan struct{}) - go func(handler RelationHandler) { - defer func() { - if r := recover(); r != nil { - klog.Errorf("Relation handler panicked: %v, treating as failed", r) + klog.Errorf("%s panicked: %v, treating as failed", panicMsg, r) } close(done) }() diff --git a/resourcetopo/node_topology_updater.go b/resourcetopo/node_topology_updater.go index d4a37f2..f146225 100644 --- a/resourcetopo/node_topology_updater.go +++ b/resourcetopo/node_topology_updater.go @@ -69,7 +69,7 @@ func (s *nodeStorage) OnUpdate(oldObj, newObj interface{}) { } var resolvedLabelRelations []ResourceRelation - resolvedDirectRelations := make(map[nodeMeta]interface{}) + resolvedDirectRelations := make(map[nodeMeta]struct{}) for _, resolver := range s.resolvers { relations := resolver.Resolve(newTopoObj) for _, relation := range relations { @@ -83,7 +83,7 @@ func (s *nodeStorage) OnUpdate(oldObj, newObj interface{}) { cluster: relation.Cluster, namespace: directRef.Namespace, name: directRef.Name, - }] = nil + }] = struct{}{} } } } @@ -220,7 +220,7 @@ func (s *nodeStorage) addLabelResourceRelation(node *nodeInfo, relation *Resourc } } -func (s *nodeStorage) setDirectResourceRelation(node *nodeInfo, directRefs map[nodeMeta]interface{}) { +func (s *nodeStorage) setDirectResourceRelation(node *nodeInfo, directRefs map[nodeMeta]struct{}) { var toDel []*nodeInfo rangeNodeList(node.directReferredPostOrders, func(postNode *nodeInfo) { postNodeMeta := nodeMeta{ From 04bdb9127e2630f7765e6bf1cdfcdf1914436dc8 Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Sat, 25 Jul 2026 23:38:38 +0800 Subject: [PATCH 8/9] feat(resourcetopo): add metrics for timeouts, panics, active goroutines, queue depth - Add execMetrics struct with timeouts/panics/active counters passed to the generic executeHandlers function. - Add Metrics() to Manager interface returning Metrics struct with: NodeHandlerTimeouts, RelationHandlerTimeouts, NodeHandlerPanics, RelationHandlerPanics, ActiveHandlerGoroutines, NodeEventQueueDepth, RelationEventQueueDepth. - Counters are atomic.Int64, updated inside executeHandlers. - Tests verify timeout, panic, and active goroutine metrics. Co-Authored-By: Claude --- resourcetopo/event_queue.go | 40 ++++++++----- resourcetopo/manager.go | 28 ++++++++- resourcetopo/resourcetopo_test.go | 96 +++++++++++++++++++++++++++++++ resourcetopo/types.go | 4 ++ 4 files changed, 152 insertions(+), 16 deletions(-) diff --git a/resourcetopo/event_queue.go b/resourcetopo/event_queue.go index aa1d797..f0d74fc 100644 --- a/resourcetopo/event_queue.go +++ b/resourcetopo/event_queue.go @@ -19,6 +19,7 @@ package resourcetopo import ( "context" "strings" + "sync/atomic" "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -77,35 +78,36 @@ func (m *manager) handleNodeEvent() { } ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) + nodeMetrics := &execMetrics{ + timeouts: &m.nodeHandlerTimeouts, + panics: &m.nodeHandlerPanics, + active: &m.activeHandlerGoroutines, + } switch evtType { case EventTypeAdd: executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnAdd(ctx, node) }, "Node handler", func() { - m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnAdd timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - }) + }, nodeMetrics) case EventTypeUpdate: executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnUpdate(ctx, node) }, "Node handler", func() { - m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - }) + }, nodeMetrics) case EventTypeDelete: executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnDelete(ctx, node) }, "Node handler", func() { - m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnDelete timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - }) + }, nodeMetrics) case EventTypeRelatedUpdate: executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { h.OnRelatedUpdate(ctx, node) }, "Node handler", func() { - m.nodeHandlerTimeouts.Add(1) klog.Warningf("Node handler OnRelatedUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) - }) + }, nodeMetrics) } cancel() m.nodeEventQueue.Done(item) @@ -144,43 +146,54 @@ func (m *manager) handleRelationEvent() { } ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) + relationMetrics := &execMetrics{ + timeouts: &m.relationHandlerTimeouts, + panics: &m.relationHandlerPanics, + active: &m.activeHandlerGoroutines, + } switch evtType { case EventTypeAdd: executeHandlers(ctx, handlers, func(h RelationHandler) { h.OnAdd(ctx, preNode, postNode) }, "Relation handler", func() { - m.relationHandlerTimeouts.Add(1) klog.Warningf("Relation handler OnAdd timeout for relation %s/%s -> %s/%s, skipping remaining handlers", preNode.namespace, preNode.name, postNode.namespace, postNode.name) - }) + }, relationMetrics) case EventTypeDelete: executeHandlers(ctx, handlers, func(h RelationHandler) { h.OnDelete(ctx, preNode, postNode) }, "Relation handler", func() { - m.relationHandlerTimeouts.Add(1) klog.Warningf("Relation handler OnDelete timeout for relation %s/%s -> %s/%s, skipping remaining handlers", preNode.namespace, preNode.name, postNode.namespace, postNode.name) - }) + }, relationMetrics) } cancel() m.relationEventQueue.Done(item) } } +// execMetrics holds atomic counters updated by executeHandlers. +type execMetrics struct { + timeouts, panics, active *atomic.Int64 +} + // executeHandlers executes handlers with timeout support using generics. // Each handler runs in its own goroutine. If the context deadline fires before // a handler completes, remaining handlers are skipped and onTimeout is called. -func executeHandlers[H any](ctx context.Context, handlers []H, exec func(H), panicMsg string, onTimeout func()) { +func executeHandlers[H any](ctx context.Context, handlers []H, exec func(H), panicMsg string, onTimeout func(), m *execMetrics) { timeout := false for _, h := range handlers { if timeout { break } done := make(chan struct{}) + m.active.Add(1) go func(handler H) { defer func() { + m.active.Add(-1) if r := recover(); r != nil { klog.Errorf("%s panicked: %v, treating as failed", panicMsg, r) + m.panics.Add(1) } close(done) }() @@ -191,6 +204,7 @@ func executeHandlers[H any](ctx context.Context, handlers []H, exec func(H), pan // Handler completed case <-ctx.Done(): timeout = true + m.timeouts.Add(1) onTimeout() } } diff --git a/resourcetopo/manager.go b/resourcetopo/manager.go index d66015b..e8fbb43 100644 --- a/resourcetopo/manager.go +++ b/resourcetopo/manager.go @@ -45,11 +45,33 @@ type manager struct { // metrics nodeHandlerTimeouts atomic.Int64 relationHandlerTimeouts atomic.Int64 + nodeHandlerPanics atomic.Int64 + relationHandlerPanics atomic.Int64 + activeHandlerGoroutines atomic.Int64 } -// HandlerTimeoutMetrics returns the count of timed-out node and relation handlers. -func (m *manager) HandlerTimeoutMetrics() (nodeTimeouts, relationTimeouts int64) { - return m.nodeHandlerTimeouts.Load(), m.relationHandlerTimeouts.Load() +// Metrics returns runtime counters for monitoring. +type Metrics struct { + NodeHandlerTimeouts int64 + RelationHandlerTimeouts int64 + NodeHandlerPanics int64 + RelationHandlerPanics int64 + ActiveHandlerGoroutines int64 + NodeEventQueueDepth int + RelationEventQueueDepth int +} + +// Metrics returns current handler and queue counters. +func (m *manager) Metrics() Metrics { + return Metrics{ + NodeHandlerTimeouts: m.nodeHandlerTimeouts.Load(), + RelationHandlerTimeouts: m.relationHandlerTimeouts.Load(), + NodeHandlerPanics: m.nodeHandlerPanics.Load(), + RelationHandlerPanics: m.relationHandlerPanics.Load(), + ActiveHandlerGoroutines: m.activeHandlerGoroutines.Load(), + NodeEventQueueDepth: m.nodeEventQueue.Len(), + RelationEventQueueDepth: m.relationEventQueue.Len(), + } } func NewResourcesTopoManager(cfg ManagerConfig) (Manager, error) { diff --git a/resourcetopo/resourcetopo_test.go b/resourcetopo/resourcetopo_test.go index 35d872e..5296812 100644 --- a/resourcetopo/resourcetopo_test.go +++ b/resourcetopo/resourcetopo_test.go @@ -2532,3 +2532,99 @@ func (h *panickingNodeHandler) OnDelete(ctx context.Context, info NodeInfo) func (h *panickingNodeHandler) OnRelatedUpdate(ctx context.Context, info NodeInfo) { panic("intentional panic in OnRelatedUpdate") } var _ NodeHandler = &panickingNodeHandler{} + +var _ = Describe("test handler metrics", func() { + var manager Manager + var fakeClient *fake.Clientset + var ctx context.Context + var cancel func() + + BeforeEach(func() { + fakeClient = fake.NewSimpleClientset() + k8sInformerFactory := informers.NewSharedInformerFactory(fakeClient, 0) + var err error + cfg := buildManagerConfig(buildClusterTest(k8sInformerFactory)) + cfg.EventHandlerTimeout = 100 * time.Millisecond + manager, err = NewResourcesTopoManager(*cfg) + Expect(err).NotTo(HaveOccurred()) + + ctx, cancel = context.WithCancel(context.Background()) + manager.Start(ctx.Done()) + k8sInformerFactory.Start(ctx.Done()) + k8sInformerFactory.WaitForCacheSync(ctx.Done()) + }) + AfterEach(func() { + cancel() + time.Sleep(200 * time.Millisecond) + }) + + It("timeout metric is incremented when handler times out", func() { + crName := "crMetricsTimeout" + + blockCh := make(chan struct{}) + blockingHandler := &blockingNodeHandler{blockCh: blockCh} + Expect(manager.AddNodeHandler(ClusterRoleMeta, blockingHandler)).To(Succeed()) + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + + // Trigger timeout + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + time.Sleep(500 * time.Millisecond) + + metrics := manager.Metrics() + Expect(metrics.NodeHandlerTimeouts).To(BeNumerically(">", 0), + fmt.Sprintf("Expected node handler timeout count > 0, got %d", metrics.NodeHandlerTimeouts)) + }) + + It("panic metric is incremented when handler panics", func() { + crName := "crMetricsPanic" + + panickingHandler := &panickingNodeHandler{} + Expect(manager.AddNodeHandler(ClusterRoleMeta, panickingHandler)).To(Succeed()) + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + + metricsBefore := manager.Metrics() + + // Trigger panic + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + time.Sleep(500 * time.Millisecond) + + metricsAfter := manager.Metrics() + Expect(metricsAfter.NodeHandlerPanics).To(BeNumerically(">", metricsBefore.NodeHandlerPanics), + fmt.Sprintf("Expected panic count to increase: before=%d, after=%d", + metricsBefore.NodeHandlerPanics, metricsAfter.NodeHandlerPanics)) + }) + + It("active goroutine count returns to zero after handlers complete", func() { + crName := "crMetricsActive" + + blockCh := make(chan struct{}) + blockingHandler := &blockingNodeHandler{blockCh: blockCh} + Expect(manager.AddNodeHandler(ClusterRoleMeta, blockingHandler)).To(Succeed()) + saHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ServiceAccountMeta, saHandler)).To(Succeed()) + crbHandler := &objecthandler{} + Expect(manager.AddNodeHandler(ClusterRoleBindingMeta, crbHandler)).To(Succeed()) + + // Trigger handler — cooperative handler returns on ctx.Done() + Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) + time.Sleep(500 * time.Millisecond) + + metrics := manager.Metrics() + Expect(metrics.ActiveHandlerGoroutines).To(BeNumerically("==", 0), + fmt.Sprintf("Expected 0 active goroutines, got %d", metrics.ActiveHandlerGoroutines)) + }) + + It("queue depth metrics are available", func() { + // Queue should be empty initially after stabilization + time.Sleep(100 * time.Millisecond) + metrics := manager.Metrics() + Expect(metrics.NodeEventQueueDepth).To(BeNumerically(">=", 0)) + Expect(metrics.RelationEventQueueDepth).To(BeNumerically(">=", 0)) + }) +}) diff --git a/resourcetopo/types.go b/resourcetopo/types.go index db2b7a8..741e9f7 100644 --- a/resourcetopo/types.go +++ b/resourcetopo/types.go @@ -78,6 +78,10 @@ type Manager interface { // GetNode return the ref to Node that match the resouorce meta and node's name if existed. GetNode(meta metav1.TypeMeta, name types.NamespacedName) (NodeInfo, error) + + // Metrics returns runtime counters for monitoring handler timeouts, panics, + // active goroutines, and event queue depths. + Metrics() Metrics } type TopoNodeStorage interface { From 61cdbdce4509afa4b52384659344cf5f9b85db44 Mon Sep 17 00:00:00 2001 From: "yuyinglu.yyl" Date: Fri, 31 Jul 2026 15:36:05 +0800 Subject: [PATCH 9/9] =?UTF-8?q?fix(resourcetopo):=20fix=20lint=20issues=20?= =?UTF-8?q?=E2=80=94=20misspell=20'cancelled'=E2=86=92'canceled',=20shadow?= =?UTF-8?q?=20builtin=20'min',=20and=20gofmt=20formatting?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude --- resourcetopo/doc.go | 81 +++++++++++++++---------------- resourcetopo/resourcetopo_test.go | 39 ++++++++++----- resourcetopo/types.go | 14 +++--- 3 files changed, 74 insertions(+), 60 deletions(-) diff --git a/resourcetopo/doc.go b/resourcetopo/doc.go index fb10e30..1594a0a 100644 --- a/resourcetopo/doc.go +++ b/resourcetopo/doc.go @@ -18,7 +18,7 @@ // resources. It maintains relationships between different resource types and notifies // handlers when resources or their relationships change. // -// Lock Hierarchy +// # Lock Hierarchy // // To prevent deadlocks, locks must be acquired in the following order: // @@ -50,55 +50,54 @@ // - Protects: nodeUpdateHandler, relationUpdateHandler // - Independent - can be acquired anytime // -// Threading Model +// # Threading Model // // This package uses multiple goroutines for concurrent processing. Understanding the // threading model is essential for avoiding race conditions and deadlocks. // // 1. Informer Threads (External): -// - Source: Kubernetes informers (one per resource type) -// - Entry points: nodeStorage.OnAdd, OnUpdate, OnDelete -// - Purpose: Receives resource change events from API server -// - Lock behavior: Acquires metaLock.Lock() during updateNodeMeta(), relationsLock during relation updates -// - Note: These are the ONLY threads that write to metaLock (labels, ownerRefs, objectExisted) +// - Source: Kubernetes informers (one per resource type) +// - Entry points: nodeStorage.OnAdd, OnUpdate, OnDelete +// - Purpose: Receives resource change events from API server +// - Lock behavior: Acquires metaLock.Lock() during updateNodeMeta(), relationsLock during relation updates +// - Note: These are the ONLY threads that write to metaLock (labels, ownerRefs, objectExisted) // // 2. Event Processor Threads (Internal): -// - Source: Created by Start() -> startHandleEvent() -// - Count: 2 goroutines -// - handleNodeEvent(): Processes node add/update/delete/relatedUpdate events -// - handleRelationEvent(): Processes relation add/delete events -// - Trigger: newNodeEvent(), newRelationEvent() queue events via workqueue -// - Lock behavior: Read-only access to handlers via RLock, no direct node locking -// - Note: Reads handlers while holding no locks - relies on handlersLock for registration +// - Source: Created by Start() -> startHandleEvent() +// - Count: 2 goroutines +// - handleNodeEvent(): Processes node add/update/delete/relatedUpdate events +// - handleRelationEvent(): Processes relation add/delete events +// - Trigger: newNodeEvent(), newRelationEvent() queue events via workqueue +// - Lock behavior: Read-only access to handlers via RLock, no direct node locking +// - Note: Reads handlers while holding no locks - relies on handlersLock for registration // // 3. User/Client Threads (External): -// - Source: User code calling Manager methods -// - Examples: GetNode(), GetTopoNodeStorage(), AddNodeHandler() -// - Lock behavior: Uses storagesLock.RLock() for reads -// - Note: These are typically short-lived operations +// - Source: User code calling Manager methods +// - Examples: GetNode(), GetTopoNodeStorage(), AddNodeHandler() +// - Lock behavior: Uses storagesLock.RLock() for reads +// - Note: These are typically short-lived operations // // 4. Callback Threads (External - via Event Processors): -// - Source: User-provided NodeHandler and RelationHandler callbacks -// - Entry: Called by handleNodeEvent/handleRelationEvent -// - Execution: Each callback runs in its own goroutine with a configurable timeout -// (ManagerConfig.EventHandlerTimeout, default 1s). If a callback does not return -// within the timeout, a warning is logged and all remaining handlers for that -// event are skipped. -// - Lock behavior: No resourcetopo locks held during callback execution -// - Concurrency: Multiple callbacks (across different handlers or events) may run -// concurrently. Handler implementations that share mutable state must be -// thread-safe. -// - Safety: Handlers receive node references but should not cache them long-term -// -// Thread Safety Guidelines -// -// - NodeInfo references returned by GetNode() are safe for concurrent reads -// - Do NOT call AddTopologyConfig() after Start() - not thread safe -// - Handler registration (AddNodeHandler, AddRelationHandler) is thread safe -// - Callbacks (OnAdd, OnUpdate, etc.) run in goroutines with a timeout. -// Do not block indefinitely — timed-out callbacks cause remaining handlers -// to be skipped. Implementations must be thread-safe if they share state. -// - nodeInfo.lock can be held for extended periods during relation changes -// - metaLock is never held during callbacks or cross-node operations -// +// - Source: User-provided NodeHandler and RelationHandler callbacks +// - Entry: Called by handleNodeEvent/handleRelationEvent +// - Execution: Each callback runs in its own goroutine with a configurable timeout +// (ManagerConfig.EventHandlerTimeout, default 1s). If a callback does not return +// within the timeout, a warning is logged and all remaining handlers for that +// event are skipped. +// - Lock behavior: No resourcetopo locks held during callback execution +// - Concurrency: Multiple callbacks (across different handlers or events) may run +// concurrently. Handler implementations that share mutable state must be +// thread-safe. +// - Safety: Handlers receive node references but should not cache them long-term +// +// # Thread Safety Guidelines +// +// - NodeInfo references returned by GetNode() are safe for concurrent reads +// - Do NOT call AddTopologyConfig() after Start() - not thread safe +// - Handler registration (AddNodeHandler, AddRelationHandler) is thread safe +// - Callbacks (OnAdd, OnUpdate, etc.) run in goroutines with a timeout. +// Do not block indefinitely — timed-out callbacks cause remaining handlers +// to be skipped. Implementations must be thread-safe if they share state. +// - nodeInfo.lock can be held for extended periods during relation changes +// - metaLock is never held during callbacks or cross-node operations package resourcetopo diff --git a/resourcetopo/resourcetopo_test.go b/resourcetopo/resourcetopo_test.go index 5296812..3b780d0 100644 --- a/resourcetopo/resourcetopo_test.go +++ b/resourcetopo/resourcetopo_test.go @@ -19,6 +19,7 @@ package resourcetopo import ( "context" "fmt" + "runtime" "sync" "testing" "time" @@ -30,7 +31,6 @@ import ( "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes/fake" - "runtime" "k8s.io/klog/v2" ) @@ -2023,8 +2023,8 @@ var _ = Describe("test suite for multi routine", func() { g.Expect(len(rsNode.GetPostOrders())).To(Equal(totalNum / len(rsNames))) } }, "10s").Should(Succeed()) - }) }) +}) var _ = Describe("test handler timeout protection", func() { var manager Manager @@ -2127,18 +2127,21 @@ func (h *blockingNodeHandler) OnAdd(ctx context.Context, info NodeInfo) { case <-ctx.Done(): } } + func (h *blockingNodeHandler) OnUpdate(ctx context.Context, info NodeInfo) { select { case <-h.blockCh: case <-ctx.Done(): } } + func (h *blockingNodeHandler) OnDelete(ctx context.Context, info NodeInfo) { select { case <-h.blockCh: case <-ctx.Done(): } } + func (h *blockingNodeHandler) OnRelatedUpdate(ctx context.Context, info NodeInfo) { select { case <-h.blockCh: @@ -2159,6 +2162,7 @@ func (h *blockingRelationHandler) OnAdd(ctx context.Context, preNode, postNode N case <-ctx.Done(): } } + func (h *blockingRelationHandler) OnDelete(ctx context.Context, preNode, postNode NodeInfo) { select { case <-h.blockCh: @@ -2441,7 +2445,7 @@ var _ = Describe("test goroutine leak and panic handling", func() { crName := "crNoLeak" // Register a cooperative blocking handler for ClusterRole - // (selects on ctx.Done() so it returns when context is cancelled) + // (selects on ctx.Done() so it returns when context is canceled) blockCh := make(chan struct{}) blockingHandler := &blockingNodeHandler{blockCh: blockCh} Expect(manager.AddNodeHandler(ClusterRoleMeta, blockingHandler)).To(Succeed()) @@ -2458,7 +2462,7 @@ var _ = Describe("test goroutine leak and panic handling", func() { // Count baseline goroutines baseline := countStableGoroutines(3, 100*time.Millisecond) - // Trigger handler — ctx will be cancelled after timeout, handler returns + // Trigger handler — ctx will be canceled after timeout, handler returns Expect(fakeClient.RbacV1().ClusterRoles().Create(ctx, newClusterRole(crName), metav1.CreateOptions{})).NotTo(BeNil()) // Wait well past timeout for handler to return and goroutine to exit @@ -2512,24 +2516,35 @@ var _ = Describe("test goroutine leak and panic handling", func() { // countStableGoroutines samples NumGoroutine multiple times and returns // the minimum observed count (to filter out transient goroutines). func countStableGoroutines(samples int, interval time.Duration) int { - min := runtime.NumGoroutine() + minGoroutines := runtime.NumGoroutine() for i := 1; i < samples; i++ { time.Sleep(interval) n := runtime.NumGoroutine() - if n < min { - min = n + if n < minGoroutines { + minGoroutines = n } } - return min + return minGoroutines } // panickingNodeHandler panics on every callback, used to test panic absorption. type panickingNodeHandler struct{} -func (h *panickingNodeHandler) OnAdd(ctx context.Context, info NodeInfo) { panic("intentional panic in OnAdd") } -func (h *panickingNodeHandler) OnUpdate(ctx context.Context, info NodeInfo) { panic("intentional panic in OnUpdate") } -func (h *panickingNodeHandler) OnDelete(ctx context.Context, info NodeInfo) { panic("intentional panic in OnDelete") } -func (h *panickingNodeHandler) OnRelatedUpdate(ctx context.Context, info NodeInfo) { panic("intentional panic in OnRelatedUpdate") } +func (h *panickingNodeHandler) OnAdd(ctx context.Context, info NodeInfo) { + panic("intentional panic in OnAdd") +} + +func (h *panickingNodeHandler) OnUpdate(ctx context.Context, info NodeInfo) { + panic("intentional panic in OnUpdate") +} + +func (h *panickingNodeHandler) OnDelete(ctx context.Context, info NodeInfo) { + panic("intentional panic in OnDelete") +} + +func (h *panickingNodeHandler) OnRelatedUpdate(ctx context.Context, info NodeInfo) { + panic("intentional panic in OnRelatedUpdate") +} var _ NodeHandler = &panickingNodeHandler{} diff --git a/resourcetopo/types.go b/resourcetopo/types.go index 741e9f7..66beb36 100644 --- a/resourcetopo/types.go +++ b/resourcetopo/types.go @@ -52,7 +52,7 @@ type ManagerConfig struct { RelationEventHandleRateMaxDelay time.Duration // EventHandlerTimeout is the timeout for each event handler callback. - // If a handler takes longer than this timeout, it will be cancelled and + // If a handler takes longer than this timeout, it will be canceled and // remaining handlers will be skipped. A warning will be logged. // Default: 1 second EventHandlerTimeout time.Duration @@ -176,30 +176,30 @@ type ResourceRelation struct { type RelationHandler interface { // OnAdd is called after relation newly added between preOrder and postOrder. - // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + // The ctx is canceled when the handler exceeds EventHandlerTimeout. OnAdd(ctx context.Context, preOrder, postOrder NodeInfo) // OnDelete is called after relation deleted between preOrder and postOrder. - // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + // The ctx is canceled when the handler exceeds EventHandlerTimeout. OnDelete(ctx context.Context, preOrder, postOrder NodeInfo) } type NodeHandler interface { // OnAdd is called when the node related resource object is added. - // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + // The ctx is canceled when the handler exceeds EventHandlerTimeout. OnAdd(ctx context.Context, info NodeInfo) // OnUpdate is called when the node related resource object is updated. - // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + // The ctx is canceled when the handler exceeds EventHandlerTimeout. OnUpdate(ctx context.Context, info NodeInfo) // OnDelete is called when the node related resource object is deleted. - // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + // The ctx is canceled when the handler exceeds EventHandlerTimeout. OnDelete(ctx context.Context, info NodeInfo) // OnRelatedUpdate is called when the nodes' // post-order(or pre-order if ReverseNotice configured) object is updated. - // The ctx is cancelled when the handler exceeds EventHandlerTimeout. + // The ctx is canceled when the handler exceeds EventHandlerTimeout. OnRelatedUpdate(ctx context.Context, info NodeInfo) }