diff --git a/resourcetopo/doc.go b/resourcetopo/doc.go new file mode 100644 index 00000000..1594a0ac --- /dev/null +++ b/resourcetopo/doc.go @@ -0,0 +1,103 @@ +/** + * 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.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 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 +// - 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/event_queue.go b/resourcetopo/event_queue.go index 8d653d61..f0d74fcb 100644 --- a/resourcetopo/event_queue.go +++ b/resourcetopo/event_queue.go @@ -17,7 +17,9 @@ package resourcetopo import ( + "context" "strings" + "sync/atomic" "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -41,6 +43,7 @@ const ( defaultRelationEventHandleRateMaxDelay = time.Second defaultNodeEventHandlePeriod = time.Second defaultRelationEventHandlePeriod = time.Second + defaultEventHandlerTimeout = time.Second ) func (m *manager) startHandleEvent(stopCh <-chan struct{}) { @@ -51,43 +54,62 @@ 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 } + ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) + nodeMetrics := &execMetrics{ + timeouts: &m.nodeHandlerTimeouts, + panics: &m.nodeHandlerPanics, + active: &m.activeHandlerGoroutines, + } switch evtType { case EventTypeAdd: - for _, h := range storage.nodeUpdateHandler { - h.OnAdd(node) - } + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnAdd(ctx, node) + }, "Node handler", func() { + klog.Warningf("Node handler OnAdd timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }, nodeMetrics) case EventTypeUpdate: - for _, h := range storage.nodeUpdateHandler { - h.OnUpdate(node) - } + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnUpdate(ctx, node) + }, "Node handler", func() { + klog.Warningf("Node handler OnUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }, nodeMetrics) case EventTypeDelete: - for _, h := range storage.nodeUpdateHandler { - h.OnDelete(node) - } + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnDelete(ctx, node) + }, "Node handler", func() { + klog.Warningf("Node handler OnDelete timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }, nodeMetrics) case EventTypeRelatedUpdate: - for _, h := range storage.nodeUpdateHandler { - h.OnRelatedUpdate(node) - } + executeHandlers(ctx, storage.nodeUpdateHandler, func(h NodeHandler) { + h.OnRelatedUpdate(ctx, node) + }, "Node handler", func() { + klog.Warningf("Node handler OnRelatedUpdate timeout for node %s/%s, skipping remaining handlers", node.namespace, node.name) + }, nodeMetrics) } + cancel() m.nodeEventQueue.Done(item) } } @@ -102,39 +124,95 @@ 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 } + ctx, cancel := context.WithTimeout(context.Background(), m.eventHandlerTimeout) + relationMetrics := &execMetrics{ + timeouts: &m.relationHandlerTimeouts, + panics: &m.relationHandlerPanics, + active: &m.activeHandlerGoroutines, + } switch evtType { case EventTypeAdd: - for _, handler := range handlers { - handler.OnAdd(preNode, postNode) - } + executeHandlers(ctx, handlers, func(h RelationHandler) { + h.OnAdd(ctx, preNode, postNode) + }, "Relation handler", func() { + 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: - for _, handler := range handlers { - handler.OnDelete(preNode, postNode) - } + executeHandlers(ctx, handlers, func(h RelationHandler) { + h.OnDelete(ctx, preNode, postNode) + }, "Relation handler", func() { + 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(), 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) + }() + exec(handler) + }(h) + select { + case <-done: + // Handler completed + case <-ctx.Done(): + timeout = true + m.timeouts.Add(1) + 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) m.nodeEventQueue.AddRateLimited(key) } diff --git a/resourcetopo/examples/base_example.go b/resourcetopo/examples/base_example.go index ca093eda..4bfa6da5 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 d9ddacc0..e8fbb436 100644 --- a/resourcetopo/manager.go +++ b/resourcetopo/manager.go @@ -20,6 +20,8 @@ import ( "container/list" "fmt" "sync" + "sync/atomic" + "time" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" @@ -35,7 +37,41 @@ 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 + + // metrics + nodeHandlerTimeouts atomic.Int64 + relationHandlerTimeouts atomic.Int64 + nodeHandlerPanics atomic.Int64 + relationHandlerPanics atomic.Int64 + activeHandlerGoroutines atomic.Int64 +} + +// 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) { @@ -44,9 +80,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 +139,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,8 +151,10 @@ 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 { - s := m.storages[generateMetaKey(meta)] - if s == nil { + m.storagesLock.RLock() + s, ok := m.storages[generateMetaKey(meta)] + m.storagesLock.RUnlock() + if !ok || s == nil { return fmt.Errorf("resource %v not configured in this manager", meta) } s.addNodeHandler(handler) @@ -135,18 +176,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, ok := m.storages[generateMetaKey(meta)] + m.storagesLock.RUnlock() + if !ok || 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, ok := m.storages[generateMetaKey(meta)] + m.storagesLock.RUnlock() + if !ok || s == nil { + return nil, fmt.Errorf("resource %v not configured in this manager", meta) } return s.GetNode(name) } @@ -157,7 +202,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 +215,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 +249,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 +284,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 4e40c438..50ce7a94 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.RLock() + objectExisted := n.objectExisted + n.metaLock.RUnlock() + + // 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 75a1221e..ec3fc376 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/node_storage.go b/resourcetopo/node_storage.go index a153436a..dad2c525 100644 --- a/resourcetopo/node_storage.go +++ b/resourcetopo/node_storage.go @@ -166,17 +166,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 +216,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 } diff --git a/resourcetopo/node_topology_updater.go b/resourcetopo/node_topology_updater.go index d4a37f25..f1462251 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{ diff --git a/resourcetopo/resourcetopo_test.go b/resourcetopo/resourcetopo_test.go index 85f2d92f..3b780d08 100644 --- a/resourcetopo/resourcetopo_test.go +++ b/resourcetopo/resourcetopo_test.go @@ -19,6 +19,7 @@ package resourcetopo import ( "context" "fmt" + "runtime" "sync" "testing" "time" @@ -2024,3 +2025,621 @@ var _ = Describe("test suite for multi routine", func() { }, "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 canceled) + 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 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 + 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 { + minGoroutines := runtime.NumGoroutine() + for i := 1; i < samples; i++ { + time.Sleep(interval) + n := runtime.NumGoroutine() + if n < minGoroutines { + minGoroutines = n + } + } + 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") +} + +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/resourcetopo_utils_test.go b/resourcetopo/resourcetopo_utils_test.go index d09b6ea2..d3b411c5 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 af8741f5..66beb36d 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" @@ -49,6 +50,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 canceled and + // remaining handlers will be skipped. A warning will be logged. + // Default: 1 second + EventHandlerTimeout time.Duration } type Manager interface { @@ -71,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 { @@ -164,26 +175,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 canceled 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 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 - OnAdd(info NodeInfo) + // OnAdd is called when the node related resource object is added. + // 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 - OnUpdate(info NodeInfo) + // OnUpdate is called when the node related resource object is updated. + // 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 - OnDelete(info NodeInfo) + // OnDelete is called when the node related resource object is deleted. + // 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 - OnRelatedUpdate(info NodeInfo) + // post-order(or pre-order if ReverseNotice configured) object is updated. + // The ctx is canceled when the handler exceeds EventHandlerTimeout. + OnRelatedUpdate(ctx context.Context, info NodeInfo) } // Object define get namespace/name/labels for now,