-
Notifications
You must be signed in to change notification settings - Fork 15
/
notify.go
62 lines (51 loc) · 2 KB
/
notify.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
/*
Copyright 2020 YANDEX LLC
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 hasql
// updateSubscriber represents a waiter for newly checked node event
type updateSubscriber[T Querier] struct {
ch chan *Node[T]
criterion NodeStateCriterion
}
// addUpdateSubscriber adds new dubscriber to notification pool
func (cl *Cluster[T]) addUpdateSubscriber(criterion NodeStateCriterion) <-chan *Node[T] {
// buffered channel is essential
// read WaitForNode function for more information
ch := make(chan *Node[T], 1)
cl.subscribersMu.Lock()
defer cl.subscribersMu.Unlock()
cl.subscribers = append(cl.subscribers, updateSubscriber[T]{ch: ch, criterion: criterion})
return ch
}
// notifyUpdateSubscribers sends appropriate nodes to registered subsribers.
// This function uses newly checked nodes to avoid race conditions
func (cl *Cluster[T]) notifyUpdateSubscribers(nodes CheckedNodes[T]) {
cl.subscribersMu.Lock()
defer cl.subscribersMu.Unlock()
if len(cl.subscribers) == 0 {
return
}
var nodelessWaiters []updateSubscriber[T]
// Notify all waiters
for _, subscriber := range cl.subscribers {
node := pickNodeByCriterion(nodes, cl.picker, subscriber.criterion)
if node == nil {
// Put waiter back
nodelessWaiters = append(nodelessWaiters, subscriber)
continue
}
// We won't block here, read addUpdateWaiter function for more information
subscriber.ch <- node
// No need to close channel since we write only once and forget it so does the 'client'
}
cl.subscribers = nodelessWaiters
}