Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Fix: CacheListener to use thread-safe ListenerSet for managing listeners #2769

Merged
Merged
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
78 changes: 61 additions & 17 deletions metadata/report/zookeeper/listener.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,52 @@ import (
"dubbo.apache.org/dubbo-go/v3/remoting/zookeeper"
)

// ListenerSet defines a thread-safe set of listeners
type ListenerSet struct {
sync.RWMutex
listeners map[mapping.MappingListener]struct{}
}

func NewListenerSet() *ListenerSet {
return &ListenerSet{
listeners: make(map[mapping.MappingListener]struct{}),
}
}

// Add adds a listener to the set
func (s *ListenerSet) Add(listener mapping.MappingListener) {
s.Lock()
defer s.Unlock()
s.listeners[listener] = struct{}{}
}

// Remove removes a listener from the set
func (s *ListenerSet) Remove(listener mapping.MappingListener) {
s.Lock()
defer s.Unlock()
delete(s.listeners, listener)
}

// Has checks if a listener exists in the set
func (s *ListenerSet) Has(listener mapping.MappingListener) bool {
s.RLock()
defer s.RUnlock()
_, ok := s.listeners[listener]
return ok
}

// ForEach iterates over all listeners in the set
func (s *ListenerSet) ForEach(f func(mapping.MappingListener) error) error {
s.RLock()
defer s.RUnlock()
for listener := range s.listeners {
if err := f(listener); err != nil {
return err
}
}
return nil
}

// CacheListener defines keyListeners and rootPath
type CacheListener struct {
// key is zkNode Path and value is set of listeners
Expand All @@ -57,35 +103,33 @@ func (l *CacheListener) AddListener(key string, listener mapping.MappingListener
if err != nil {
return
}
listeners, loaded := l.keyListeners.LoadOrStore(key, map[mapping.MappingListener]struct{}{listener: {}})
if loaded {
listeners.(map[mapping.MappingListener]struct{})[listener] = struct{}{}
l.keyListeners.Store(key, listeners)
}
// try to store the new set. If key exists, add listener to existing set
listeners, _ := l.keyListeners.LoadOrStore(key, NewListenerSet())
listeners.(*ListenerSet).Add(listener)
}

// RemoveListener will delete a listener if loaded
func (l *CacheListener) RemoveListener(key string, listener mapping.MappingListener) {
listeners, loaded := l.keyListeners.Load(key)
if loaded {
delete(listeners.(map[mapping.MappingListener]struct{}), listener)
listeners.(*ListenerSet).Remove(listener)
}
}

// DataChange changes all listeners' event
func (l *CacheListener) DataChange(event remoting.Event) bool {
if listeners, ok := l.keyListeners.Load(event.Path); ok {
for listener := range listeners.(map[mapping.MappingListener]struct{}) {
appNames := strings.Split(event.Content, constant.CommaSeparator)
set := gxset.NewSet()
for _, e := range appNames {
set.Add(e)
}
err := listener.OnEvent(registry.NewServiceMappingChangedEvent(l.pathToKey(event.Path), set))
if err != nil {
logger.Error("Error notify mapping change event.", err)
return false
}
appNames := strings.Split(event.Content, constant.CommaSeparator)
set := gxset.NewSet()
for _, e := range appNames {
set.Add(e)
}
err := listeners.(*ListenerSet).ForEach(func(listener mapping.MappingListener) error {
return listener.OnEvent(registry.NewServiceMappingChangedEvent(l.pathToKey(event.Path), set))
})
if err != nil {
logger.Error("Error notify mapping change event.", err)
return false
}
return true
}
Expand Down
Loading