/
pubsub.go
77 lines (64 loc) · 1.82 KB
/
pubsub.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
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
// Copyright (c) 2017 Opsidian Ltd.
//
// This Source Code Form is subject to the terms of the Mozilla Public
// License, v. 2.0. If a copy of the MPL was not distributed with this
// file, You can obtain one at http://mozilla.org/MPL/2.0/.
package conflow
import (
"fmt"
"sync"
)
type subscription struct {
container *NodeContainer
next *subscription
}
type PubSub struct {
subs map[ID]*subscription
mu *sync.RWMutex
}
func NewPubSub() *PubSub {
return &PubSub{
subs: make(map[ID]*subscription),
mu: &sync.RWMutex{},
}
}
// Subscribe will subscribe the given node container for the given dependency
func (p *PubSub) Subscribe(c *NodeContainer, id ID) {
p.mu.Lock()
p.subs[id] = &subscription{container: c, next: p.subs[id]}
p.mu.Unlock()
}
// Unsubscribe will unsubscribe the given node container for the given dependency
func (p *PubSub) Unsubscribe(c *NodeContainer, id ID) {
p.mu.Lock()
if p.subs[id].container.Node().ID() == c.Node().ID() {
p.subs[id] = p.subs[id].next
p.mu.Unlock()
return
}
for sub := p.subs[id]; sub.next != nil; sub = sub.next {
if sub.next.container.Node().ID() == c.Node().ID() {
sub.next = sub.next.next
p.mu.Unlock()
return
}
}
p.mu.Unlock()
panic(fmt.Errorf("unsubscribe unsuccessful, %q was never subscribed for %q", c.Node().ID(), id))
}
// Publish will notify all node containers which are subscribed for the dependency
// The ready function will run on any containers which have all dependencies satisfied
func (p *PubSub) Publish(c Container) {
p.mu.RLock()
defer p.mu.RUnlock()
for sub := p.subs[c.Node().ID()]; sub != nil; sub = sub.next {
sub.container.SetDependency(c)
}
}
// HasSubscribers will return true if the given block has subscribers
func (p *PubSub) HasSubscribers(id ID) bool {
p.mu.RLock()
_, ok := p.subs[id]
p.mu.RUnlock()
return ok
}