|
|
@@ -16,9 +16,7 @@ package api
|
|
|
|
|
|
import (
|
|
|
"sync"
|
|
|
- "time"
|
|
|
|
|
|
- "github.com/coreos/etcd/etcdserver"
|
|
|
"github.com/coreos/etcd/version"
|
|
|
"github.com/coreos/go-semver/semver"
|
|
|
"github.com/coreos/pkg/capnslog"
|
|
|
@@ -43,45 +41,32 @@ var (
|
|
|
"3.0.0": {AuthCapability: true, V3rpcCapability: true},
|
|
|
}
|
|
|
|
|
|
- // capLoopOnce ensures we only create one capability monitor goroutine
|
|
|
- capLoopOnce sync.Once
|
|
|
-
|
|
|
enableMapMu sync.RWMutex
|
|
|
// enabledMap points to a map in capabilityMaps
|
|
|
enabledMap map[Capability]bool
|
|
|
+
|
|
|
+ curVersion *semver.Version
|
|
|
)
|
|
|
|
|
|
func init() {
|
|
|
enabledMap = make(map[Capability]bool)
|
|
|
}
|
|
|
|
|
|
-// RunCapabilityLoop checks the cluster version every 500ms and updates
|
|
|
-// the enabledMap when the cluster version increased.
|
|
|
-func RunCapabilityLoop(s *etcdserver.EtcdServer) {
|
|
|
- go capLoopOnce.Do(func() { runCapabilityLoop(s) })
|
|
|
-}
|
|
|
-
|
|
|
-func runCapabilityLoop(s *etcdserver.EtcdServer) {
|
|
|
- stopped := s.StopNotify()
|
|
|
-
|
|
|
- var pv *semver.Version
|
|
|
- for {
|
|
|
- if v := s.ClusterVersion(); v != pv {
|
|
|
- if pv == nil || (v != nil && pv.LessThan(*v)) {
|
|
|
- pv = v
|
|
|
- enableMapMu.Lock()
|
|
|
- enabledMap = capabilityMaps[pv.String()]
|
|
|
- enableMapMu.Unlock()
|
|
|
- plog.Infof("enabled capabilities for version %s", version.Cluster(pv.String()))
|
|
|
- }
|
|
|
- }
|
|
|
-
|
|
|
- select {
|
|
|
- case <-stopped:
|
|
|
- return
|
|
|
- case <-time.After(500 * time.Millisecond):
|
|
|
- }
|
|
|
+// UpdateCapability updates the enabledMap when the cluster version increases.
|
|
|
+func UpdateCapability(v *semver.Version) {
|
|
|
+ if v == nil {
|
|
|
+ // if recovered but version was never set by cluster
|
|
|
+ return
|
|
|
+ }
|
|
|
+ enableMapMu.Lock()
|
|
|
+ if curVersion != nil && !curVersion.LessThan(*v) {
|
|
|
+ enableMapMu.Unlock()
|
|
|
+ return
|
|
|
}
|
|
|
+ curVersion = v
|
|
|
+ enabledMap = capabilityMaps[curVersion.String()]
|
|
|
+ enableMapMu.Unlock()
|
|
|
+ plog.Infof("enabled capabilities for version %s", version.Cluster(v.String()))
|
|
|
}
|
|
|
|
|
|
func IsCapabilityEnabled(c Capability) bool {
|