forked from flannel-io/flannel
-
Notifications
You must be signed in to change notification settings - Fork 0
/
registry.go
114 lines (96 loc) · 2.86 KB
/
registry.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
package subnet
import (
"path"
"sync"
"time"
"github.com/coreos/flannel/Godeps/_workspace/src/github.com/coreos/go-etcd/etcd"
log "github.com/coreos/flannel/Godeps/_workspace/src/github.com/golang/glog"
)
type subnetRegistry interface {
getConfig() (*etcd.Response, error)
getSubnets() (*etcd.Response, error)
createSubnet(sn, data string, ttl uint64) (*etcd.Response, error)
updateSubnet(sn, data string, ttl uint64) (*etcd.Response, error)
watchSubnets(since uint64, stop chan bool) (*etcd.Response, error)
}
type etcdSubnetRegistry struct {
mux sync.Mutex
cli *etcd.Client
endpoint string
prefix string
}
func newEtcdSubnetRegistry(endpoint, prefix string) subnetRegistry {
return &etcdSubnetRegistry{
cli: etcd.NewClient([]string{endpoint}),
endpoint: endpoint,
prefix: prefix,
}
}
func (esr *etcdSubnetRegistry) getConfig() (*etcd.Response, error) {
key := path.Join(esr.prefix, "config")
resp, err := esr.client().Get(key, false, false)
if err != nil {
return nil, err
}
return resp, nil
}
func (esr *etcdSubnetRegistry) getSubnets() (*etcd.Response, error) {
key := path.Join(esr.prefix, "subnets")
return esr.client().Get(key, false, true)
}
func (esr *etcdSubnetRegistry) createSubnet(sn, data string, ttl uint64) (*etcd.Response, error) {
key := path.Join(esr.prefix, "subnets", sn)
resp, err := esr.client().Create(key, data, ttl)
if err != nil {
return nil, err
}
ensureExpiration(resp, ttl)
return resp, nil
}
func (esr *etcdSubnetRegistry) updateSubnet(sn, data string, ttl uint64) (*etcd.Response, error) {
key := path.Join(esr.prefix, "subnets", sn)
resp, err := esr.client().Set(key, data, ttl)
if err != nil {
return nil, err
}
ensureExpiration(resp, ttl)
return resp, nil
}
func (esr *etcdSubnetRegistry) watchSubnets(since uint64, stop chan bool) (*etcd.Response, error) {
for {
key := path.Join(esr.prefix, "subnets")
resp, err := esr.client().RawWatch(key, since, true, nil, stop)
if err != nil {
if err == etcd.ErrWatchStoppedByUser {
return nil, nil
} else {
return nil, err
}
}
if len(resp.Body) == 0 {
// etcd timed out, go back but recreate the client as the underlying
// http transport gets hosed (http://code.google.com/p/go/issues/detail?id=8648)
esr.resetClient()
continue
}
return resp.Unmarshal()
}
}
func (esr *etcdSubnetRegistry) client() *etcd.Client {
esr.mux.Lock()
defer esr.mux.Unlock()
return esr.cli
}
func (esr *etcdSubnetRegistry) resetClient() {
esr.mux.Lock()
defer esr.mux.Unlock()
esr.cli = etcd.NewClient([]string{esr.endpoint})
}
func ensureExpiration(resp *etcd.Response, ttl uint64) {
if resp.Node.Expiration == nil {
// should not be but calc it ourselves in this case
log.Info("Expiration field missing on etcd response, calculating locally")
exp := time.Now().Add(time.Duration(ttl) * time.Second)
resp.Node.Expiration = &exp
}
}