forked from harness/gitness
-
Notifications
You must be signed in to change notification settings - Fork 0
/
pubsub.go
104 lines (87 loc) · 2.73 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
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
package pubsub
import (
"sync"
"code.google.com/p/go.net/context"
)
type PubSub struct {
sync.Mutex
// In-memory list of all channels being managed by the broker.
channels map[interface{}]*Channel
}
// NewPubSub creates a new instance of the PubSub type
// and returns a pointer.
func NewPubSub() *PubSub {
return &PubSub{
channels: make(map[interface{}]*Channel),
}
}
// Lookup performs a thread safe operation to return a pointer
// to an existing Channel object with the given key. If the
// Channel does not exist a nil value is returned.
func (b *PubSub) Lookup(key interface{}) *Channel {
b.Lock()
defer b.Unlock()
// find the channel in the existing list
return b.channels[key]
}
// Register performs a thread safe operation to return a pointer
// to a Channel object for the given key. The Channel is created
// if it does not yet exist.
func (b *PubSub) Register(key interface{}) *Channel {
return b.RegisterOpts(key, DefaultOpts)
}
// Register performs a thread safe operation to return a pointer
// to a Channel object for the given key. The Channel is created
// if it does not yet exist using custom options.
func (b *PubSub) RegisterOpts(key interface{}, opts *Opts) *Channel {
b.Lock()
defer b.Unlock()
// find the channel in the existing list
c, ok := b.channels[key]
if ok {
return c
}
// create the channel and register
// with the pubsub server
c = NewChannel(opts)
b.channels[key] = c
go c.start()
return c
}
// Unregister performs a thread safe operation to delete the
// Channel with the given key.
func (b *PubSub) Unregister(key interface{}) {
b.Lock()
defer b.Unlock()
// find the channel in the existing list
c, ok := b.channels[key]
if !ok {
return
}
c.Close()
delete(b.channels, key)
return
}
// Lookup performs a thread safe operation to return a pointer
// to an existing Channel object with the given key. If the
// Channel does not exist a nil value is returned.
func Lookup(c context.Context, key interface{}) *Channel {
return FromContext(c).Lookup(key)
}
// Register performs a thread safe operation to return a pointer
// to a Channel object for the given key. The Channel is created
// if it does not yet exist.
func Register(c context.Context, key interface{}) *Channel {
return FromContext(c).Register(key)
}
// Register performs a thread safe operation to return a pointer
// to a Channel object for the given key. The Channel is created
// if it does not yet exist using custom options.
func RegisterOpts(c context.Context, key interface{}, opts *Opts) *Channel {
return FromContext(c).RegisterOpts(key, opts)
}
// Unregister performs a thread safe operation to delete the
// Channel with the given key.
func Unregister(c context.Context, key interface{}) {
FromContext(c).Unregister(key)
}