forked from st3v/go-plugins
/
memory.go
128 lines (107 loc) · 2.49 KB
/
memory.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
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
// Package memory provides an in-memory registry
package memory
import (
"sync"
"time"
"github.com/micro/go-micro/cmd"
"github.com/micro/go-micro/registry"
"github.com/pborman/uuid"
"golang.org/x/net/context"
)
type memoryRegistry struct {
sync.RWMutex
services map[string][]*registry.Service
watchers map[string]*memoryWatcher
}
var (
timeout = time.Millisecond * 10
)
func init() {
cmd.DefaultRegistries["memory"] = NewRegistry
}
func (m *memoryRegistry) watch(r *registry.Result) {
var watchers []*memoryWatcher
m.RLock()
for _, w := range m.watchers {
watchers = append(watchers, w)
}
m.RUnlock()
for _, w := range watchers {
select {
case <-w.exit:
m.Lock()
delete(m.watchers, w.id)
m.Unlock()
default:
select {
case w.res <- r:
case <-time.After(timeout):
}
}
}
}
func (m *memoryRegistry) GetService(service string) ([]*registry.Service, error) {
m.RLock()
s, ok := m.services[service]
if !ok || len(s) == 0 {
m.RUnlock()
return nil, registry.ErrNotFound
}
m.RUnlock()
return s, nil
}
func (m *memoryRegistry) ListServices() ([]*registry.Service, error) {
m.RLock()
var services []*registry.Service
for _, service := range m.services {
services = append(services, service...)
}
m.RUnlock()
return services, nil
}
func (m *memoryRegistry) Register(s *registry.Service, opts ...registry.RegisterOption) error {
go m.watch(®istry.Result{Action: "update", Service: s})
m.Lock()
services := addServices(m.services[s.Name], []*registry.Service{s})
m.services[s.Name] = services
m.Unlock()
return nil
}
func (m *memoryRegistry) Deregister(s *registry.Service) error {
go m.watch(®istry.Result{Action: "delete", Service: s})
m.Lock()
services := delServices(m.services[s.Name], []*registry.Service{s})
m.services[s.Name] = services
m.Unlock()
return nil
}
func (m *memoryRegistry) Watch() (registry.Watcher, error) {
w := &memoryWatcher{
exit: make(chan bool),
res: make(chan *registry.Result),
id: uuid.NewUUID().String(),
}
m.Lock()
m.watchers[w.id] = w
m.Unlock()
return w, nil
}
func (m *memoryRegistry) String() string {
return "memory"
}
func NewRegistry(opts ...registry.Option) registry.Registry {
options := registry.Options{
Context: context.Background(),
}
for _, o := range opts {
o(&options)
}
services := getServices(options.Context)
if services == nil {
services = make(map[string][]*registry.Service)
}
return &memoryRegistry{
services: services,
watchers: make(map[string]*memoryWatcher),
}
}