From 487b7e189b2f3187525048ed0b7bca281084159a Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 12 Mar 2020 18:15:30 +0800 Subject: [PATCH 1/7] cloudcommon: options: add constants for lock method --- pkg/cloudcommon/database.go | 4 ++-- pkg/cloudcommon/options/options.go | 5 +++++ 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/pkg/cloudcommon/database.go b/pkg/cloudcommon/database.go index 061f550649..722842637a 100644 --- a/pkg/cloudcommon/database.go +++ b/pkg/cloudcommon/database.go @@ -53,11 +53,11 @@ func InitDB(options *common_options.DBOptions) { sqlchemy.SetDB(dbConn) switch options.LockmanMethod { - case "inmemory", "": + case common_options.LockMethodInMemory, "": log.Infof("using inmemory lockman") lm := lockman.NewInMemoryLockManager() lockman.Init(lm) - case "etcd": + case common_options.LockMethodEtcd: log.Infof("using etcd lockman") tlsCfg, err := options.GetEtcdTLSConfig() if err != nil { diff --git a/pkg/cloudcommon/options/options.go b/pkg/cloudcommon/options/options.go index bf6e3d137f..114e07dda5 100644 --- a/pkg/cloudcommon/options/options.go +++ b/pkg/cloudcommon/options/options.go @@ -89,6 +89,11 @@ type BaseOptions struct { structarg.BaseOptions } +const ( + LockMethodInMemory = "inmemory" + LockMethodEtcd = "etcd" +) + type CommonOptions struct { AuthURL string `help:"Keystone auth URL" alias:"auth-uri"` AdminUser string `help:"Admin username"` From 411bff8301c2b70d5a764f9b00ed2666571f98ef Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 12 Mar 2020 17:02:21 +0800 Subject: [PATCH 2/7] cloudcommon: remove default value for etcd lock prefix --- pkg/cloudcommon/db/lockman/etcd.go | 2 +- pkg/cloudcommon/options/options.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/pkg/cloudcommon/db/lockman/etcd.go b/pkg/cloudcommon/db/lockman/etcd.go index 34d43707e9..64462421da 100644 --- a/pkg/cloudcommon/db/lockman/etcd.go +++ b/pkg/cloudcommon/db/lockman/etcd.go @@ -155,7 +155,7 @@ func (config *SEtcdLockManagerConfig) validate() error { config.LockPrefix = strings.TrimSpace(config.LockPrefix) config.LockPrefix = strings.TrimRight(config.LockPrefix, "/") if config.LockPrefix == "" { - config.LockPrefix = "/" + return fmt.Errorf("empty etcd lock prefix") } return nil diff --git a/pkg/cloudcommon/options/options.go b/pkg/cloudcommon/options/options.go index 114e07dda5..5e5d1ccb34 100644 --- a/pkg/cloudcommon/options/options.go +++ b/pkg/cloudcommon/options/options.go @@ -122,7 +122,7 @@ type DBOptions struct { QueryOffsetOptimization bool `help:"apply query offset optimization"` LockmanMethod string `help:"method for lock synchronization" choices:"inmemory|etcd" default:"inmemory"` - EtcdLockPrefix string `help:"prefix of etcd lock records" default:"/locks"` + EtcdLockPrefix string `help:"prefix of etcd lock records"` EtcdLockTTL int `help:"ttl of etcd lock records"` EtcdEndpoints []string `help:"endpoints of etcd cluster"` From 20f3c775d7db6f30daddd2c13694c02ed0d0362e Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 12 Mar 2020 16:19:18 +0800 Subject: [PATCH 3/7] cloudcommon: reword help text for IsSlaveNode --- pkg/cloudcommon/options/options.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/cloudcommon/options/options.go b/pkg/cloudcommon/options/options.go index 5e5d1ccb34..1fe335b5fb 100644 --- a/pkg/cloudcommon/options/options.go +++ b/pkg/cloudcommon/options/options.go @@ -75,7 +75,7 @@ type BaseOptions struct { ConfigSyncPeriodSeconds int `help:"service config sync interval in seconds, default 300 seconds/5 minutes" default:"300"` - IsSlaveNode bool `help:"Region service slave node"` + IsSlaveNode bool `help:"Slave mode"` CronJobWorkerCount int `help:"Cron job worker count" default:"4"` DefaultQuotaValue string `help:"default quota value" choices:"unlimit|zero|default" default:"default"` From 48a0bb420b57a1410103f6c4814f165650e2411d Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 12 Mar 2020 17:31:13 +0800 Subject: [PATCH 4/7] cloudcommon: elect support --- pkg/cloudcommon/elect/doc.go | 1 + pkg/cloudcommon/elect/elect.go | 212 +++++++++++++++++++++++++++++++++ 2 files changed, 213 insertions(+) create mode 100644 pkg/cloudcommon/elect/doc.go create mode 100644 pkg/cloudcommon/elect/elect.go diff --git a/pkg/cloudcommon/elect/doc.go b/pkg/cloudcommon/elect/doc.go new file mode 100644 index 0000000000..03dc698e6e --- /dev/null +++ b/pkg/cloudcommon/elect/doc.go @@ -0,0 +1 @@ +package elect // import "yunion.io/x/onecloud/pkg/cloudcommon/elect" diff --git a/pkg/cloudcommon/elect/elect.go b/pkg/cloudcommon/elect/elect.go new file mode 100644 index 0000000000..24f11857f3 --- /dev/null +++ b/pkg/cloudcommon/elect/elect.go @@ -0,0 +1,212 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package elect + +import ( + "context" + "crypto/tls" + "sync" + "time" + + "go.etcd.io/etcd/clientv3" + "go.etcd.io/etcd/clientv3/concurrency" + "google.golang.org/grpc" + + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/cloudcommon/options" +) + +type EtcdConfig struct { + Endpoints []string + Username string + Password string + TLS *tls.Config + + LockTTL int + LockPrefix string + + opts *options.DBOptions +} + +func NewEtcdConfigFromDBOptions(opts *options.DBOptions) (*EtcdConfig, error) { + tlsCfg, err := opts.GetEtcdTLSConfig() + if err != nil { + return nil, errors.Wrap(err, "etcd tls config") + } + + config := &EtcdConfig{ + Endpoints: opts.EtcdEndpoints, + Username: opts.EtcdUsername, + Password: opts.EtcdPassword, + LockTTL: opts.EtcdLockTTL, + LockPrefix: opts.EtcdLockPrefix, + TLS: tlsCfg, + } + if config.LockTTL <= 0 { + config.LockTTL = 5 + } + return config, nil +} + +type Elect struct { + cli *clientv3.Client + path string + ttl int + + stopFunc context.CancelFunc + + config *EtcdConfig + + mutex *sync.Mutex + subscribers []chan ElectEvent +} + +type ticket struct { + session *concurrency.Session + mutex *concurrency.Mutex +} + +func (t *ticket) tearup(ctx context.Context) { + if t.session != nil { + t.session.Close() + t.session = nil + } +} + +func NewElect(config *EtcdConfig, key string) (*Elect, error) { + cli, err := clientv3.New(clientv3.Config{ + Endpoints: config.Endpoints, + Username: config.Username, + Password: config.Password, + TLS: config.TLS, + + DialOptions: []grpc.DialOption{ + grpc.WithBlock(), + grpc.WithTimeout(500 * time.Millisecond), + }, + DialTimeout: 3 * time.Second, + }) + if err != nil { + return nil, errors.Wrap(err, "new etcd client") + } + elect := &Elect{ + cli: cli, + path: config.LockPrefix + "/" + key, + ttl: config.LockTTL, + mutex: &sync.Mutex{}, + } + return elect, nil +} + +func (elect *Elect) Stop() { + elect.stopFunc() +} + +func (elect *Elect) Start(ctx context.Context) { + ctx, elect.stopFunc = context.WithCancel(ctx) + + prev := ElectEventLost + for { + select { + case <-ctx.Done(): + log.Infof("elect bye") + return + default: + now := ElectEventWin + ticket, err := elect.do(ctx) + if err != nil { + ticket.tearup(ctx) + now = ElectEventLost + log.Errorf("elect error: %v", err) + } + if now != prev { + log.Infof("notify elect event: %s -> %s", prev, now) + prev = now + elect.notify(ctx, now) + } + if err == nil { + select { + case <-ctx.Done(): + ticket.tearup(ctx) + case <-ticket.session.Done(): + } + } else { + time.Sleep(3 * time.Second) + } + } + } +} + +// do joins an election. the 1st return argument ticket must always be non-nil +func (elect *Elect) do(ctx context.Context) (*ticket, error) { + r := &ticket{} + + sess, err := concurrency.NewSession( + elect.cli, + concurrency.WithTTL(elect.ttl), + ) + if err != nil { + return r, err + } + r.session = sess + + em := concurrency.NewMutex(sess, elect.path) + if err := em.Lock(ctx); err != nil { + return r, err + } + r.mutex = em + + return r, err +} + +func (elect *Elect) Subscribe(ch chan ElectEvent) { + elect.mutex.Lock() + defer elect.mutex.Unlock() + elect.subscribers = append(elect.subscribers, ch) +} + +func (elect *Elect) notify(ctx context.Context, ev ElectEvent) { + elect.mutex.Lock() + defer elect.mutex.Unlock() + + for _, ch := range elect.subscribers { + select { + case ch <- ev: + case <-ctx.Done(): + return + default: + } + } +} + +type ElectEvent int + +const ( + ElectEventWin ElectEvent = iota + ElectEventLost +) + +func (ev ElectEvent) String() string { + switch ev { + case ElectEventWin: + return "win" + case ElectEventLost: + return "lost" + default: + return "unexpected" + } +} From 1d8fa1a1f258400805c4560bcc63a9253642161e Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 12 Mar 2020 18:42:08 +0800 Subject: [PATCH 5/7] cronman: use context for start, stop --- pkg/cloudcommon/cronman/cronman.go | 20 ++++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index b7a5aa801a..312db04239 100644 --- a/pkg/cloudcommon/cronman/cronman.go +++ b/pkg/cloudcommon/cronman/cronman.go @@ -120,7 +120,7 @@ func (cjth *CronJobTimerHeap) Pop() interface{} { type SCronJobManager struct { jobs CronJobTimerHeap - stop chan struct{} + stopFunc context.CancelFunc add chan struct{} running bool workers *appsrv.SWorkerManager @@ -134,7 +134,6 @@ func InitCronJobManager(isDbWorker bool, workerCount int) *SCronJobManager { workers: appsrv.NewWorkerManager("CronJobWorkers", workerCount, 1024, isDbWorker), dataLock: new(sync.Mutex), add: make(chan struct{}), - stop: make(chan struct{}), } } return manager @@ -299,6 +298,12 @@ func (self *SCronJobManager) next(now time.Time) { } func (self *SCronJobManager) Start() { + ctx := context.Background() + ctx, self.stopFunc = context.WithCancel(ctx) + self.start(ctx) +} + +func (self *SCronJobManager) start(ctx context.Context) { if self.running { return } @@ -306,13 +311,11 @@ func (self *SCronJobManager) Start() { defer self.dataLock.Unlock() self.running = true self.init() - go self.run() + go self.run(ctx) } func (self *SCronJobManager) Stop() { - if self.stop != nil { - close(self.stop) - } + self.stopFunc() } func (self *SCronJobManager) init() { @@ -327,7 +330,7 @@ func (self *SCronJobManager) init() { } } -func (self *SCronJobManager) run() { +func (self *SCronJobManager) run(ctx context.Context) { var timer *time.Timer var now = time.Now() for { @@ -343,8 +346,9 @@ func (self *SCronJobManager) run() { self.runJobs(now) case <-self.add: continue - case <-self.stop: + case <-ctx.Done(): timer.Stop() + self.running = false return } } From 0f7ff65a465e5a138465fa463941d8022bf60e0a Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 12 Mar 2020 18:53:46 +0800 Subject: [PATCH 6/7] cronman: add Start2() for working with electObj --- pkg/cloudcommon/cronman/cronman.go | 28 ++++++++++++++++++++++++++++ 1 file changed, 28 insertions(+) diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index 312db04239..569cc48855 100644 --- a/pkg/cloudcommon/cronman/cronman.go +++ b/pkg/cloudcommon/cronman/cronman.go @@ -27,6 +27,7 @@ import ( "yunion.io/x/onecloud/pkg/appctx" "yunion.io/x/onecloud/pkg/appsrv" + "yunion.io/x/onecloud/pkg/cloudcommon/elect" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/auth" ) @@ -297,6 +298,33 @@ func (self *SCronJobManager) next(now time.Time) { } } +func (self *SCronJobManager) Start2(ctx context.Context, electObj *elect.Elect) { + ctx, self.stopFunc = context.WithCancel(ctx) + if electObj == nil { + self.start(ctx) + return + } + + go func() { + ch := make(chan elect.ElectEvent) + electObj.Subscribe(ch) + for { + select { + case ev := <-ch: + log.Infof("cronman: elect event %s: cronman", ev) + switch ev { + case elect.ElectEventWin: + self.start(ctx) + case elect.ElectEventLost: + self.Stop() + } + case <-ctx.Done(): + return + } + } + }() +} + func (self *SCronJobManager) Start() { ctx := context.Background() ctx, self.stopFunc = context.WithCancel(ctx) From be5322d5306b60bc915037ef25c8eb1f8629705a Mon Sep 17 00:00:00 2001 From: Yousong Zhou Date: Thu, 12 Mar 2020 18:54:15 +0800 Subject: [PATCH 7/7] region: start cronman with electObj --- pkg/compute/service/service.go | 23 +++++++++++++++++++++-- 1 file changed, 21 insertions(+), 2 deletions(-) diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index 0241dc381d..e46a86a298 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -15,6 +15,7 @@ package service import ( + "context" "os" "time" @@ -27,6 +28,7 @@ import ( app_common "yunion.io/x/onecloud/pkg/cloudcommon/app" "yunion.io/x/onecloud/pkg/cloudcommon/cronman" "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/elect" common_options "yunion.io/x/onecloud/pkg/cloudcommon/options" _ "yunion.io/x/onecloud/pkg/compute/guestdrivers" _ "yunion.io/x/onecloud/pkg/compute/hostdrivers" @@ -72,6 +74,24 @@ func StartService() { models.InitSyncWorkers(options.Options.CloudSyncWorkerCount) + var ( + electObj *elect.Elect + ctx, cancelFunc = context.WithCancel(context.Background()) + ) + defer cancelFunc() + + if opts.LockmanMethod == common_options.LockMethodEtcd { + cfg, err := elect.NewEtcdConfigFromDBOptions(dbOpts) + if err != nil { + log.Fatalf("etcd config for elect: %v", err) + } + electObj, err = elect.NewElect(cfg, "@master-role") + if err != nil { + log.Fatalf("new elect instance: %v", err) + } + go electObj.Start(ctx) + } + if !opts.IsSlaveNode { cron := cronman.InitCronJobManager(true, options.Options.CronJobWorkerCount) cron.AddJobAtIntervals("CleanPendingDeleteServers", time.Duration(opts.PendingDeleteCheckSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers) @@ -98,8 +118,7 @@ func StartService() { cron.AddJobEveryFewDays("SyncElasticCacheSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncElasticCacheSkus, true) cron.AddJobEveryFewDays("StorageSnapshotsRecycle", 1, 2, 0, 0, models.StorageManager.StorageSnapshotsRecycle, false) - cron.Start() - defer cron.Stop() + go cron.Start2(ctx, electObj) } app_common.ServeForever(app, baseOpts)