This is an automated email from the ASF dual-hosted git repository.

alexstocks pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/dubbo-go.git


The following commit(s) were added to refs/heads/develop by this push:
     new cade2982d fix: improve thread safety in CacheListener (#2769)
cade2982d is described below

commit cade2982d6e51f3e0826d8cf32402e4957986667
Author: solisamicus <[email protected]>
AuthorDate: Tue Dec 24 13:56:04 2024 +0800

    fix: improve thread safety in CacheListener (#2769)
---
 metadata/report/zookeeper/listener.go | 78 +++++++++++++++++++++++++++--------
 1 file changed, 61 insertions(+), 17 deletions(-)

diff --git a/metadata/report/zookeeper/listener.go 
b/metadata/report/zookeeper/listener.go
index 84e79f6a4..e844f11c1 100644
--- a/metadata/report/zookeeper/listener.go
+++ b/metadata/report/zookeeper/listener.go
@@ -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
@@ -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
        }

Reply via email to