Skip to content
Open
103 changes: 103 additions & 0 deletions resourcetopo/doc.go
Original file line number Diff line number Diff line change
@@ -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
114 changes: 96 additions & 18 deletions resourcetopo/event_queue.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@
package resourcetopo

import (
"context"
"strings"
"sync/atomic"
"time"

metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
Expand All @@ -41,6 +43,7 @@ const (
defaultRelationEventHandleRateMaxDelay = time.Second
defaultNodeEventHandlePeriod = time.Second
defaultRelationEventHandlePeriod = time.Second
defaultEventHandlerTimeout = time.Second
)

func (m *manager) startHandleEvent(stopCh <-chan struct{}) {
Expand All @@ -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)
}
}
Expand All @@ -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)
}

Expand Down
12 changes: 6 additions & 6 deletions resourcetopo/examples/base_example.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()))
Expand All @@ -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())
}

Expand Down
Loading
Loading