diff --git a/whisper/whisperv5/filter.go b/whisper/whisperv5/filter.go index d8c558db26..4f7a0861f4 100644 --- a/whisper/whisperv5/filter.go +++ b/whisper/whisperv5/filter.go @@ -27,7 +27,7 @@ import ( ) const ( - ALL_TOPICS = "" + ALL_TOPICS = "" MAX_POOL_CAPACITY = 1000 ) @@ -45,17 +45,17 @@ type Filter struct { } type Filters struct { - watchers map[string]*Filter - whisper *Whisper - mutex sync.RWMutex - topicMatcher *topicMatcher + watchers map[string]*Filter + whisper *Whisper + mutex sync.RWMutex + topicMatcher *topicMatcher } func NewFilters(w *Whisper) *Filters { fs := &Filters{ - watchers: make(map[string]*Filter), - whisper: w, - topicMatcher: newTopicMatcher(), + watchers: make(map[string]*Filter), + whisper: w, + topicMatcher: newTopicMatcher(), } return fs } @@ -228,6 +228,7 @@ func IsPubKeyEqual(a, b *ecdsa.PublicKey) bool { return a.X.Cmp(b.X) == 0 && a.Y.Cmp(b.Y) == 0 } +//topicMatcher keeps topic->watcher mapping type topicMatcher struct { //structure - map[topic]map[filterID] //mapping for topics @@ -237,6 +238,7 @@ type topicMatcher struct { pool sync.Pool } +//newTopicMatcher returns a newly created topic matcher func newTopicMatcher() *topicMatcher { tm := new(topicMatcher) tm.mapper = make(map[string]map[string]struct{}) @@ -247,9 +249,12 @@ func newTopicMatcher() *topicMatcher { return tm } +//take returns []string from pool func (fs *topicMatcher) take() []string { return fs.pool.Get().([]string) } + +//resolve put []string to pool func (fs *topicMatcher) resolve(s []string) { if cap(s) > MAX_POOL_CAPACITY { return @@ -257,6 +262,7 @@ func (fs *topicMatcher) resolve(s []string) { fs.pool.Put(s[:0]) } +//addFilterToTopicsMapping fill topic->watcher mapping for current watcher func (fs *topicMatcher) addFilterToTopicsMapping(watcher *Filter, id string) { fs.mx.Lock() defer fs.mx.Unlock() @@ -271,6 +277,7 @@ func (fs *topicMatcher) addFilterToTopicsMapping(watcher *Filter, id string) { } } +//removeTopicFromTopicMapping removes mapping info by filterID func (fs *topicMatcher) removeTopicFromTopicMapping(id string) { fs.mx.Lock() defer fs.mx.Unlock() @@ -279,6 +286,7 @@ func (fs *topicMatcher) removeTopicFromTopicMapping(id string) { } } +//prepareTopicsMapping returns set of topics for watcher func (fs *topicMatcher) prepareTopicsMapping(watcher *Filter) map[string]struct{} { topics := make(map[string]struct{}, len(watcher.Topics)) @@ -294,6 +302,7 @@ func (fs *topicMatcher) prepareTopicsMapping(watcher *Filter) map[string]struct{ return topics } +//matchedTopics write all matched topics to matched func (fs *topicMatcher) matchedTopics(topic TopicType, matched *[]string) { fs.mx.RLock() defer fs.mx.RUnlock() diff --git a/whisper/whisperv6/filter.go b/whisper/whisperv6/filter.go index aca2a297c7..a59da6ddae 100644 --- a/whisper/whisperv6/filter.go +++ b/whisper/whisperv6/filter.go @@ -89,6 +89,7 @@ func (fs *Filters) Install(watcher *Filter) (string, error) { } fs.watchers[id] = watcher + //add topic matching for watcher fs.topicMatcher.addFilterToTopicsMapping(watcher, id) return id, err } @@ -246,6 +247,7 @@ func IsPubKeyEqual(a, b *ecdsa.PublicKey) bool { return a.X.Cmp(b.X) == 0 && a.Y.Cmp(b.Y) == 0 } +//topicMatcher keeps topic->watcher mapping type topicMatcher struct { //structure - map[topic]map[filterID] //mapping for topics @@ -255,6 +257,7 @@ type topicMatcher struct { pool sync.Pool } +//newTopicMatcher returns a newly created topic matcher func newTopicMatcher() *topicMatcher { tm := new(topicMatcher) tm.mapper = make(map[string]map[string]struct{}) @@ -265,9 +268,12 @@ func newTopicMatcher() *topicMatcher { return tm } +//take returns []string from pool func (fs *topicMatcher) take() []string { return fs.pool.Get().([]string) } + +//resolve put []string to pool func (fs *topicMatcher) resolve(s []string) { if cap(s) > MAX_POOL_CAPACITY { return @@ -275,6 +281,7 @@ func (fs *topicMatcher) resolve(s []string) { fs.pool.Put(s[:0]) } +//addFilterToTopicsMapping fill topic->watcher mapping for current watcher func (fs *topicMatcher) addFilterToTopicsMapping(watcher *Filter, id string) { fs.mx.Lock() defer fs.mx.Unlock() @@ -289,14 +296,16 @@ func (fs *topicMatcher) addFilterToTopicsMapping(watcher *Filter, id string) { } } -func (fs *topicMatcher) removeTopicFromTopicMapping(id string) { +//removeTopicFromTopicMapping removes mapping info by filterID +func (fs *topicMatcher) removeTopicFromTopicMapping(filterID string) { fs.mx.Lock() defer fs.mx.Unlock() for i := range fs.mapper { - delete(fs.mapper[i], id) + delete(fs.mapper[i], filterID) } } +//prepareTopicsMapping returns set of topics for watcher func (fs *topicMatcher) prepareTopicsMapping(watcher *Filter) map[string]struct{} { topics := make(map[string]struct{}, len(watcher.Topics)) @@ -312,6 +321,7 @@ func (fs *topicMatcher) prepareTopicsMapping(watcher *Filter) map[string]struct{ return topics } +//matchedTopics write all matched topics to matched func (fs *topicMatcher) matchedTopics(topic TopicType, matched *[]string) { fs.mx.RLock() defer fs.mx.RUnlock()