/
etcd_resolver.go
124 lines (111 loc) · 2.65 KB
/
etcd_resolver.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
package rpc
import (
"context"
"github.com/spf13/viper"
"github.com/zenpk/dorm-system/pkg/zap"
clientv3 "go.etcd.io/etcd/client/v3"
"google.golang.org/grpc/resolver"
"sync"
)
type EtcdResolverBuilder struct {
etcdClient *clientv3.Client
}
func InitEtcdResolverBuilder() (*EtcdResolverBuilder, error) {
etcdClient, err := initEtcdClient()
if err != nil {
return nil, err
}
rb := new(EtcdResolverBuilder)
rb.etcdClient = etcdClient
return rb, nil
}
func (r EtcdResolverBuilder) Close() error {
return r.etcdClient.Close()
}
func (r EtcdResolverBuilder) Scheme() string {
return viper.GetString("etcd.scheme")
}
func (r EtcdResolverBuilder) Build(target resolver.Target, cc resolver.ClientConn, opts resolver.BuildOptions) (resolver.Resolver, error) {
prefix := target.URL.Scheme + "://" + target.URL.Host + target.URL.Path
res, err := r.etcdClient.Get(context.Background(), prefix, clientv3.WithPrefix())
if err != nil {
return nil, err
}
ctx, cancelFunc := context.WithCancel(context.Background())
er := &etcdResolver{
etcdClient: r.etcdClient,
cc: cc,
ctx: ctx,
cancel: cancelFunc,
scheme: target.URL.Scheme,
}
for _, kv := range res.Kvs {
er.store(kv.Key, kv.Value)
}
if err := er.updateState(); err != nil {
return nil, err
}
go er.watcher()
return er, nil
}
type etcdResolver struct {
etcdClient *clientv3.Client
cc resolver.ClientConn
ctx context.Context
cancel context.CancelFunc
scheme string
ipPool sync.Map
}
func (e *etcdResolver) ResolveNow(resolver.ResolveNowOptions) {
}
func (e *etcdResolver) Close() {
e.cancel()
}
func (e *etcdResolver) watcher() {
watchChan := e.etcdClient.Watch(e.ctx, e.scheme)
for {
select {
case val := <-watchChan:
for _, event := range val.Events {
switch event.Type {
case 0:
e.store(event.Kv.Key, event.Kv.Value)
if err := e.updateState(); err != nil {
zap.Logger.Error(err)
return
}
case 1:
e.del(event.Kv.Key)
if err := e.updateState(); err != nil {
zap.Logger.Error(err)
return
}
}
}
case <-e.ctx.Done():
return
}
}
}
func (e *etcdResolver) store(k, v []byte) {
e.ipPool.Store(string(k), string(v))
}
func (e *etcdResolver) del(key []byte) {
e.ipPool.Delete(string(key))
}
func (e *etcdResolver) updateState() error {
var addrList resolver.State
e.ipPool.Range(func(k, v interface{}) bool {
name, ok := k.(string)
if !ok {
return false
}
addr, ok := v.(string)
if !ok {
return false
}
addrList.Addresses = append(addrList.Addresses, resolver.Address{Addr: addr, ServerName: name})
return true
})
return e.cc.UpdateState(addrList)
}