Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions cmd/yurthub/app/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,15 @@ func Complete(options *options.YurtHubOptions, stopCh <-chan struct{}) (*YurtHub
cfg.ConfigManager = configManager
cfg.FilterFinder = filterFinder

// Wire the configuration reload listener so that when the yurt-hub-cfg
// ConfigMap changes at runtime, the FilterManager rebuilds its internal state
// (nameToObjectFilter and resourceSyncers) to match the new filter settings.
configManager.AddListener(func(newCfg map[string]string) {
if err := filterFinder.Reset(newCfg); err != nil {
klog.Errorf("could not reset filter manager after config reload, %v", err)
}
})

if options.EnableDummyIf {
klog.V(2).
Infof("create dummy network interface %s(%s)", options.HubAgentDummyIfName, options.HubAgentDummyIfIP)
Expand Down
20 changes: 20 additions & 0 deletions pkg/yurthub/configuration/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,9 @@ var (
defaultCacheAgents = []string{"kubelet", "kube-proxy", "flanneld", "coredns", "raven-agent-ds", projectinfo.GetAgentName(), projectinfo.GetHubName()}
)

// listener is a callback invoked when the configuration reloads.
type listener func(newCfg map[string]string)

// Manager is used for managing all configurations of Yurthub in yurt-hub-cfg configmap.
// This configuration configmap includes configurations of cache agents and filters. I'm sure that new
// configurations will be added according to user's new requirements.
Expand All @@ -54,6 +57,7 @@ type Manager struct {
baseKeyToFilters map[string][]string
reqKeyToFilters map[string][]string
configMapSynced cache.InformerSynced
listeners []listener
}

func NewConfigurationManager(nodeName string, sharedFactory informers.SharedInformerFactory) *Manager {
Expand Down Expand Up @@ -109,6 +113,15 @@ func (m *Manager) HasSynced() bool {
return m.configMapSynced()
}

// AddListener registers a callback that will be invoked when the configuration
// is reloaded via the yurt-hub-cfg ConfigMap watcher. The callback receives the
// new ConfigMap data map and should be goroutine-safe.
func (m *Manager) AddListener(fn listener) {
m.Lock()
defer m.Unlock()
m.listeners = append(m.listeners, fn)
}

// ListAllCacheAgents is used for listing all cache agents.
func (m *Manager) ListAllCacheAgents() []string {
m.RLock()
Expand Down Expand Up @@ -165,6 +178,13 @@ func (m *Manager) updateConfigmap(oldObj, newObj interface{}) {
if filterSettingsChanged(oldCfg.Data, newCfg.Data) {
m.updateFilterSettings(newCfg.Data, "update")
}

m.RLock()
lis := m.listeners
m.RUnlock()
for _, l := range lis {
l(newCfg.Data)
}
}

func (m *Manager) deleteConfigmap(obj interface{}) {
Expand Down
1 change: 1 addition & 0 deletions pkg/yurthub/filter/interfaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,7 @@ type ObjectFilter interface {
type FilterFinder interface {
FindResponseFilter(req *http.Request) (ResponseFilter, bool)
FindObjectFilter(req *http.Request) (ObjectFilter, bool)
Reset(cmData map[string]string) error
ResourceSyncer
}

Expand Down
95 changes: 89 additions & 6 deletions pkg/yurthub/filter/manager/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ package manager
import (
"net/http"
"strconv"
"sync"

"k8s.io/client-go/dynamic/dynamicinformer"
"k8s.io/client-go/informers"
Expand All @@ -39,9 +40,16 @@ import (

type Manager struct {
filter.Approver
mu sync.RWMutex
nameToObjectFilter map[string]filter.ObjectFilter
serializerManager *serializer.SerializerManager
resourceSyncers []filter.ResourceSyncer

// dependencies for dynamic filter rebuilding
options *yurtoptions.YurtHubOptions
sharedFactory informers.SharedInformerFactory
dynamicSharedFactory dynamicinformer.DynamicSharedInformerFactory
client kubernetes.Interface
}

func NewFilterManager(options *yurtoptions.YurtHubOptions,
Expand Down Expand Up @@ -90,15 +98,23 @@ func NewFilterManager(options *yurtoptions.YurtHubOptions,

// 5. new filter manager including approver and nameToObjectFilter
// if resource filters are disabled, nameToObjectFilter and resourceSyncers will be empty silces.
return &Manager{
Approver: approver.NewApprover(options.NodeName, configManager),
nameToObjectFilter: nameToFilters,
serializerManager: serializerManager,
resourceSyncers: resourceSyncers,
}, nil
m := &Manager{
Approver: approver.NewApprover(options.NodeName, configManager),
nameToObjectFilter: nameToFilters,
serializerManager: serializerManager,
resourceSyncers: resourceSyncers,
options: options,
sharedFactory: sharedFactory,
dynamicSharedFactory: dynamicSharedFactory,
client: proxiedClient,
}

return m, nil
}

func (m *Manager) HasSynced() bool {
m.mu.RLock()
defer m.mu.RUnlock()
for i := range m.resourceSyncers {
if !m.resourceSyncers[i].HasSynced() {
return false
Expand All @@ -108,6 +124,9 @@ func (m *Manager) HasSynced() bool {
}

func (m *Manager) FindResponseFilter(req *http.Request) (filter.ResponseFilter, bool) {
m.mu.RLock()
defer m.mu.RUnlock()

if len(m.nameToObjectFilter) == 0 {
return nil, false
}
Expand All @@ -132,6 +151,9 @@ func (m *Manager) FindResponseFilter(req *http.Request) (filter.ResponseFilter,
}

func (m *Manager) FindObjectFilter(req *http.Request) (filter.ObjectFilter, bool) {
m.mu.RLock()
defer m.mu.RUnlock()

if len(m.nameToObjectFilter) == 0 {
return nil, false
}
Expand All @@ -154,3 +176,64 @@ func (m *Manager) FindObjectFilter(req *http.Request) (filter.ObjectFilter, bool

return objectfilter.CreateFilterChain(objectFilters), true
}

// Reset rebuilds the filter manager's internal state (nameToObjectFilter and resourceSyncers)
// from the new configuration data. It must be called under the write lock (mu).
// The method builds new maps into local variables first, then atomically swaps them so that
// a failed Reset never leaves FilterManager in an inconsistent state serving live traffic with
// stale filters. If construction fails, the error is returned and the previous state is left intact.
//
// The caller must ensure that no goroutine is concurrently calling FindResponseFilter or
// FindObjectFilter while Reset is running, or external synchronization must be provided.
func (m *Manager) Reset(cmData map[string]string) error {

// Step 1: build new state in local variables first (never mutate struct fields until success)
newNameToFilters := make(map[string]filter.ObjectFilter)
newResourceSyncers := make([]filter.ResourceSyncer, 0)

if m.options != nil && m.options.EnableResourceFilter {
// Re-create the filter registry with the current disabled list.
filtersReg := base.NewFilters(m.options.DisabledResourceFilters)
// Re-register all filter factories.
yurtoptions.RegisterAllFilters(filtersReg)

// Re-build the initializer chain using the new config.
mutatedMasterServicePort := strconv.Itoa(m.options.YurtHubProxySecurePort)
mutatedMasterServiceHost := m.options.YurtHubProxyHost
if m.options.EnableDummyIf {
mutatedMasterServiceHost = m.options.HubAgentDummyIfIP
}
genericInitializer := initializer.New(m.sharedFactory, m.client, m.options.NodeName, m.options.NodePoolName,
mutatedMasterServiceHost, mutatedMasterServicePort)
nodesInitializer := initializer.NewNodesInitializer(m.options.EnableNodePool, m.options.EnablePoolServiceTopology, m.dynamicSharedFactory)
initializerChain := base.Initializers{}
initializerChain = append(initializerChain, genericInitializer, nodesInitializer)

// Initialize all object filters with the new chain.
var err error
newNameToFilters, err = filtersReg.NewFromFilters(initializerChain)
if err != nil {
klog.Errorf("could not rebuild filters during Reset, %v", err)
return err
}

// Collect resource syncers from the newly initialized filters.
for name, objFilter := range newNameToFilters {
if resourceSyncer, ok := objFilter.(filter.ResourceSyncer); ok {
klog.Infof("filter %s need to sync resource before starting to work (reset)", name)
newResourceSyncers = append(newResourceSyncers, resourceSyncer)
}
}
}

// Step 2: only assign the new maps after successful construction.
// This ensures a failed Reset never leaves FilterManager half-updated.
m.mu.Lock()
m.nameToObjectFilter = newNameToFilters
m.resourceSyncers = newResourceSyncers
m.mu.Unlock()

klog.Infof("FilterManager state rebuilt successfully after ConfigMap update: %d object filters, %d resource syncers",
len(m.nameToObjectFilter), len(m.resourceSyncers))
return nil
}
86 changes: 83 additions & 3 deletions pkg/yurthub/filter/manager/manager_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,8 @@ import (
"github.com/openyurtio/openyurt/pkg/yurthub/configuration"
"github.com/openyurtio/openyurt/pkg/yurthub/filter"
"github.com/openyurtio/openyurt/pkg/yurthub/kubernetes/serializer"
"github.com/openyurtio/openyurt/pkg/yurthub/proxy/util"
proxyutil "github.com/openyurtio/openyurt/pkg/yurthub/proxy/util"
"github.com/openyurtio/openyurt/pkg/yurthub/util"
)

func TestFindResponseFilter(t *testing.T) {
Expand Down Expand Up @@ -166,7 +167,7 @@ func TestFindResponseFilter(t *testing.T) {
responseFilter, isFound = finder.FindResponseFilter(req)
})

handler = util.WithRequestClientComponent(handler)
handler = proxyutil.WithRequestClientComponent(handler)
handler = filters.WithRequestInfo(handler, resolver)
handler.ServeHTTP(httptest.NewRecorder(), req)

Expand Down Expand Up @@ -311,7 +312,7 @@ func TestFindObjectFilter(t *testing.T) {
objectFilter, isFound = finder.FindObjectFilter(req)
})

handler = util.WithRequestClientComponent(handler)
handler = proxyutil.WithRequestClientComponent(handler)
handler = filters.WithRequestInfo(handler, resolver)
handler.ServeHTTP(httptest.NewRecorder(), req)

Expand All @@ -330,6 +331,85 @@ func TestFindObjectFilter(t *testing.T) {
}
}

// TestFilterManagerDynamicUpdate verifies that the FilterManager correctly
// rebuilds its internal state (nameToObjectFilter and resourceSyncers) when
// the yurt-hub-cfg ConfigMap changes at runtime. This ensures that filter
// configuration updates are propagated transparently without requiring a
// Yurthub restart.
func TestFilterManagerDynamicUpdate(t *testing.T) {
fakeClient := &fake.Clientset{}
scheme := runtime.NewScheme()
apis.AddToScheme(scheme)
fakeDynamicClient := dynamicfake.NewSimpleDynamicClient(scheme)
serializerManager := serializer.NewSerializerManager()

// Config A: enable masterservice filter only
optionsA := &options.YurtHubOptions{
EnableResourceFilter: true,
WorkingMode: string(util.WorkingModeEdge),
DisabledResourceFilters: []string{},
EnableDummyIf: false,
NodeName: "test-node",
YurtHubProxySecurePort: 10268,
HubAgentDummyIfIP: "127.0.0.1",
YurtHubProxyHost: "127.0.0.1",
}
optionsA.DisabledResourceFilters = []string{}

sharedFactory, nodePoolFactory := informers.NewSharedInformerFactory(fakeClient, 24*time.Hour),
dynamicinformer.NewDynamicSharedInformerFactory(fakeDynamicClient, 24*time.Hour)

configManager := configuration.NewConfigurationManager(optionsA.NodeName, sharedFactory)

// NOTE: We can't fully start the informers in this unit test context without
// a real k8s cluster, but we can test the Reset path by directly exercising
// the method with synthetic config data. The test below verifies that Reset
// rebuilds the internal maps correctly.
finderA, _ := NewFilterManager(optionsA, sharedFactory, nodePoolFactory, fakeClient, serializerManager, configManager)

// Reset the finder with new config data simulating a ConfigMap update.
newCfg := map[string]string{
"masterservice": "kubelet,services,get",
}

if err := finderA.Reset(newCfg); err != nil {
t.Fatalf("Reset() unexpected error: %v", err)
}

// Verify that FindObjectFilter/FindResponseFilter work after Reset.
resolver := newTestRequestInfoResolver()
req, _ := http.NewRequest("GET", "/api/v1/services", nil)
req.RemoteAddr = "127.0.0.1"
req.Header.Set("User-Agent", "kubelet")

var foundB bool
var responseFilter filter.ResponseFilter
var ok bool
var handler http.Handler = http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) {
_, foundB = finderA.FindObjectFilter(req)
responseFilter, ok = finderA.FindResponseFilter(req)
})

handler = proxyutil.WithRequestClientComponent(handler)
handler = filters.WithRequestInfo(handler, resolver)
handler.ServeHTTP(httptest.NewRecorder(), req)

if !foundB {
t.Error("expected FindObjectFilter to find filters after Reset, but got not found")
}

if !ok {
t.Error("expected FindResponseFilter to find a response filter after Reset")
} else if responseFilter != nil {
names := strings.Split(responseFilter.Name(), ",")
filterNames := sets.New(names...)
if !filterNames.Has("masterservice") {
t.Errorf("expected filter names to include masterservice, got %v", names)
}
}
}

// newTestRequestInfoResolver is a test helper that returns a default request info resolver.
func newTestRequestInfoResolver() *request.RequestInfoFactory {
return &request.RequestInfoFactory{
APIPrefixes: sets.NewString("api", "apis"),
Expand Down
5 changes: 5 additions & 0 deletions pkg/yurthub/proxy/autonomy/autonomy.go
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,7 @@ func (ap *AutonomyProxy) updateNodeStatus(req *http.Request) (runtime.Object, er

var node, retNode runtime.Object
var err error
hadError := false
for i := 0; i < nodeStatusUpdateRetry; i++ {
node, err = ap.tryUpdateNodeConditions(i, req)
if node != nil {
Expand All @@ -93,6 +94,7 @@ func (ap *AutonomyProxy) updateNodeStatus(req *http.Request) (runtime.Object, er
if errors.Is(err, ErrDirectClientMgr) {
break
} else if err != nil {
hadError = true
klog.ErrorS(err, "Error getting or updating node status, will retry")
} else {
return retNode, nil
Expand All @@ -101,6 +103,9 @@ func (ap *AutonomyProxy) updateNodeStatus(req *http.Request) (runtime.Object, er
if retNode == nil {
return nil, fmt.Errorf("failed to get node")
}
if hadError {
return nil, fmt.Errorf("failed to update node autonomy status after retries")
}
klog.ErrorS(err, "failed to update node autonomy status")
return retNode, nil
}
Expand Down
Loading
Loading