mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #5483 from yousong/feature/yousong-multi-instances
cloudcommon: elect support
This commit is contained in:
@@ -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"
|
||||
)
|
||||
@@ -120,7 +121,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 +135,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
|
||||
@@ -298,7 +298,40 @@ 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)
|
||||
self.start(ctx)
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) start(ctx context.Context) {
|
||||
if self.running {
|
||||
return
|
||||
}
|
||||
@@ -306,13 +339,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 +358,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 +374,9 @@ func (self *SCronJobManager) run() {
|
||||
self.runJobs(now)
|
||||
case <-self.add:
|
||||
continue
|
||||
case <-self.stop:
|
||||
case <-ctx.Done():
|
||||
timer.Stop()
|
||||
self.running = false
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
package elect // import "yunion.io/x/onecloud/pkg/cloudcommon/elect"
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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"`
|
||||
@@ -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"`
|
||||
@@ -117,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"`
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user