mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #5853 from wanyaoqi/feature/wyq/host-health2
host health check
This commit is contained in:
@@ -304,25 +304,7 @@ func (self *SCronJobManager) Start2(ctx context.Context, electObj *elect.Elect)
|
||||
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
|
||||
}
|
||||
}
|
||||
}()
|
||||
electObj.SubscribeWithAction(ctx, func() { self.start(ctx) }, self.Stop)
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) Start() {
|
||||
|
||||
@@ -261,6 +261,7 @@ const (
|
||||
ACT_GUEST_CREATE_FROM_IMPORT_FAIL = "guest_create_from_import_fail"
|
||||
ACT_GUEST_PANICKED = "guest_panicked"
|
||||
ACT_HOST_MAINTENANCE = "host_maintenance"
|
||||
ACT_HOST_DOWN = "host_down"
|
||||
|
||||
ACT_UPLOAD_OBJECT = "upload_obj"
|
||||
ACT_DELETE_OBJECT = "delete_obj"
|
||||
|
||||
@@ -72,7 +72,7 @@ type Elect struct {
|
||||
config *EtcdConfig
|
||||
|
||||
mutex *sync.Mutex
|
||||
subscribers []chan ElectEvent
|
||||
subscribers []chan electEvent
|
||||
}
|
||||
|
||||
type ticket struct {
|
||||
@@ -119,18 +119,18 @@ func (elect *Elect) Stop() {
|
||||
func (elect *Elect) Start(ctx context.Context) {
|
||||
ctx, elect.stopFunc = context.WithCancel(ctx)
|
||||
|
||||
prev := ElectEventLost
|
||||
prev := electEventLost
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
log.Infof("elect bye")
|
||||
return
|
||||
default:
|
||||
now := ElectEventWin
|
||||
now := electEventWin
|
||||
ticket, err := elect.do(ctx)
|
||||
if err != nil {
|
||||
ticket.tearup(ctx)
|
||||
now = ElectEventLost
|
||||
now = electEventLost
|
||||
log.Errorf("elect error: %v", err)
|
||||
}
|
||||
if now != prev {
|
||||
@@ -173,13 +173,13 @@ func (elect *Elect) do(ctx context.Context) (*ticket, error) {
|
||||
return r, err
|
||||
}
|
||||
|
||||
func (elect *Elect) Subscribe(ch chan ElectEvent) {
|
||||
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) {
|
||||
func (elect *Elect) notify(ctx context.Context, ev electEvent) {
|
||||
elect.mutex.Lock()
|
||||
defer elect.mutex.Unlock()
|
||||
|
||||
@@ -193,18 +193,49 @@ func (elect *Elect) notify(ctx context.Context, ev ElectEvent) {
|
||||
}
|
||||
}
|
||||
|
||||
type ElectEvent int
|
||||
func (elect *Elect) SubscribeWithAction(ctx context.Context, onWin, onLost func()) {
|
||||
go func() {
|
||||
ch := make(chan electEvent, 3)
|
||||
var ev electEvent
|
||||
elect.subscribe(ch)
|
||||
for {
|
||||
select {
|
||||
case ev = <-ch:
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
drain:
|
||||
for {
|
||||
select {
|
||||
case ev = <-ch:
|
||||
continue
|
||||
default:
|
||||
break drain
|
||||
}
|
||||
}
|
||||
log.Infof("elect event %s", ev)
|
||||
switch ev {
|
||||
case electEventWin:
|
||||
onWin()
|
||||
case electEventLost:
|
||||
onLost()
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
type electEvent int
|
||||
|
||||
const (
|
||||
ElectEventWin ElectEvent = iota
|
||||
ElectEventLost
|
||||
electEventWin electEvent = iota
|
||||
electEventLost
|
||||
)
|
||||
|
||||
func (ev ElectEvent) String() string {
|
||||
func (ev electEvent) String() string {
|
||||
switch ev {
|
||||
case ElectEventWin:
|
||||
case electEventWin:
|
||||
return "win"
|
||||
case ElectEventLost:
|
||||
case electEventLost:
|
||||
return "lost"
|
||||
default:
|
||||
return "unexpected"
|
||||
|
||||
@@ -22,6 +22,7 @@ import (
|
||||
"time"
|
||||
|
||||
"go.etcd.io/etcd/clientv3"
|
||||
"google.golang.org/grpc"
|
||||
|
||||
"yunion.io/x/log"
|
||||
|
||||
@@ -39,24 +40,43 @@ type SEtcdClient struct {
|
||||
|
||||
namespace string
|
||||
|
||||
leaseId clientv3.LeaseID
|
||||
leaseId clientv3.LeaseID
|
||||
onKeepaliveFailure func()
|
||||
leaseLiving bool
|
||||
|
||||
watchers map[string]*SEtcdWatcher
|
||||
}
|
||||
|
||||
func NewEtcdClient(opt *SEtcdOptions) (*SEtcdClient, error) {
|
||||
func defaultOnKeepAliveFailed() {
|
||||
log.Fatalf("etcd keepalive failed")
|
||||
}
|
||||
|
||||
func NewEtcdClient(opt *SEtcdOptions, onKeepaliveFailure func()) (*SEtcdClient, error) {
|
||||
var err error
|
||||
var tlsConfig *tls.Config
|
||||
|
||||
if opt.EtcdEnabldSsl {
|
||||
tlsConfig, err = seclib2.InitTLSConfig(opt.EtcdSslCertfile, opt.EtcdSslKeyfile)
|
||||
if err != nil {
|
||||
log.Errorf("init tls config fail %s", err)
|
||||
return nil, err
|
||||
if opt.TLSConfig == nil {
|
||||
if len(opt.EtcdSslCaCertfile) > 0 {
|
||||
tlsConfig, err = seclib2.InitTLSConfigWithCA(
|
||||
opt.EtcdSslCertfile, opt.EtcdSslKeyfile, opt.EtcdSslCaCertfile)
|
||||
} else {
|
||||
tlsConfig, err = seclib2.InitTLSConfig(opt.EtcdSslCertfile, opt.EtcdSslKeyfile)
|
||||
}
|
||||
if err != nil {
|
||||
log.Errorf("init tls config fail %s", err)
|
||||
return nil, err
|
||||
}
|
||||
} else {
|
||||
tlsConfig = opt.TLSConfig
|
||||
}
|
||||
}
|
||||
|
||||
etcdClient := &SEtcdClient{}
|
||||
if onKeepaliveFailure == nil {
|
||||
onKeepaliveFailure = defaultOnKeepAliveFailed
|
||||
}
|
||||
etcdClient.onKeepaliveFailure = onKeepaliveFailure
|
||||
|
||||
timeoutSeconds := opt.EtcdTimeoutSeconds
|
||||
if timeoutSeconds == 0 {
|
||||
@@ -69,6 +89,10 @@ func NewEtcdClient(opt *SEtcdOptions) (*SEtcdClient, error) {
|
||||
Username: opt.EtcdUsername,
|
||||
Password: opt.EtcdPassword,
|
||||
TLS: tlsConfig,
|
||||
|
||||
DialOptions: []grpc.DialOption{
|
||||
grpc.WithBlock(),
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -95,7 +119,9 @@ func NewEtcdClient(opt *SEtcdOptions) (*SEtcdClient, error) {
|
||||
|
||||
err = etcdClient.startSession()
|
||||
if err != nil {
|
||||
etcdClient.Close()
|
||||
if e := etcdClient.Close(); e != nil {
|
||||
log.Errorf("etcd client close failed %s", e)
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
return etcdClient, nil
|
||||
@@ -129,12 +155,18 @@ func (cli *SEtcdClient) startSession() error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
cli.leaseLiving = true
|
||||
|
||||
go func() {
|
||||
for {
|
||||
ka := <-ch
|
||||
if ka == nil {
|
||||
log.Fatalf("fail to keepalive")
|
||||
cli.leaseLiving = false
|
||||
log.Errorf("fail to keepalive sessoin")
|
||||
if cli.onKeepaliveFailure != nil {
|
||||
cli.onKeepaliveFailure()
|
||||
}
|
||||
break
|
||||
} else {
|
||||
log.Debugf("etcd session %d keepalive ttl: %d", ka.ID, ka.TTL)
|
||||
}
|
||||
@@ -144,6 +176,13 @@ func (cli *SEtcdClient) startSession() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (cli *SEtcdClient) RestartSession() error {
|
||||
if cli.leaseLiving {
|
||||
return errors.New("session is living, can't restart")
|
||||
}
|
||||
return cli.startSession()
|
||||
}
|
||||
|
||||
func (cli *SEtcdClient) getKey(key string) string {
|
||||
if len(cli.namespace) > 0 {
|
||||
return fmt.Sprintf("%s%s", cli.namespace, key)
|
||||
|
||||
@@ -18,13 +18,13 @@ var (
|
||||
defaultClient *SEtcdClient
|
||||
)
|
||||
|
||||
func InitDefaultEtcdClient(opt *SEtcdOptions) error {
|
||||
func InitDefaultEtcdClient(opt *SEtcdOptions, onKeepaliveFailure func()) error {
|
||||
if defaultClient != nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
var err error
|
||||
defaultClient, err = NewEtcdClient(opt)
|
||||
defaultClient, err = NewEtcdClient(opt, onKeepaliveFailure)
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
@@ -14,6 +14,8 @@
|
||||
|
||||
package etcd
|
||||
|
||||
import "crypto/tls"
|
||||
|
||||
type SEtcdOptions struct {
|
||||
EtcdEndpoint []string `help:"etcd endpoints in format of addr:port"`
|
||||
EtcdTimeoutSeconds int `default:"5" help:"etcd dial timeout in seconds"`
|
||||
@@ -25,7 +27,9 @@ type SEtcdOptions struct {
|
||||
EtcdUsername string `help:"etcd username"`
|
||||
EtcdPassword string `help:"etcd password"`
|
||||
|
||||
EtcdEnabldSsl bool `help:"enable SSL/TLS"`
|
||||
EtcdSslCertfile string `help:"ssl certification file"`
|
||||
EtcdSslKeyfile string `help:"ssl certification private key file"`
|
||||
EtcdEnabldSsl bool `help:"enable SSL/TLS"`
|
||||
EtcdSslCertfile string `help:"ssl certification file"`
|
||||
EtcdSslKeyfile string `help:"ssl certification private key file"`
|
||||
EtcdSslCaCertfile string `help:"ssl ca certification file"`
|
||||
TLSConfig *tls.Config `help:"tls config"`
|
||||
}
|
||||
|
||||
@@ -125,22 +125,25 @@ 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"`
|
||||
EtcdLockTTL int `help:"ttl of etcd lock records"`
|
||||
EtcdEndpoints []string `help:"endpoints of etcd cluster"`
|
||||
LockmanMethod string `help:"method for lock synchronization" choices:"inmemory|etcd" default:"inmemory"`
|
||||
|
||||
EtcdUsername string `help:"username of etcd cluster"`
|
||||
EtcdPassword string `help:"password of etcd cluster"`
|
||||
|
||||
EtcdUseTLS bool `help:"use tls transport to connect etcd cluster" default:"false"`
|
||||
EtcdSkipTLSVerify bool `help:"skip tls verification" default:"false"`
|
||||
EtcdCacert string `help:"path to cacert for connecting to etcd cluster"`
|
||||
EtcdCert string `help:"path to cert file for connecting to etcd cluster"`
|
||||
EtcdKey string `help:"path to key file for connecting to etcd cluster"`
|
||||
EtcdOptions
|
||||
EtcdLockPrefix string `help:"prefix of etcd lock records"`
|
||||
EtcdLockTTL int `help:"ttl of etcd lock records"`
|
||||
}
|
||||
|
||||
func (this *DBOptions) GetEtcdTLSConfig() (*tls.Config, error) {
|
||||
type EtcdOptions struct {
|
||||
EtcdEndpoints []string `help:"endpoints of etcd cluster"`
|
||||
EtcdUsername string `help:"username of etcd cluster"`
|
||||
EtcdPassword string `help:"password of etcd cluster"`
|
||||
EtcdUseTLS bool `help:"use tls transport to connect etcd cluster" default:"false"`
|
||||
EtcdSkipTLSVerify bool `help:"skip tls verification" default:"false"`
|
||||
EtcdCacert string `help:"path to cacert for connecting to etcd cluster"`
|
||||
EtcdCert string `help:"path to cert file for connecting to etcd cluster"`
|
||||
EtcdKey string `help:"path to key file for connecting to etcd cluster"`
|
||||
}
|
||||
|
||||
func (this *EtcdOptions) GetEtcdTLSConfig() (*tls.Config, error) {
|
||||
var (
|
||||
cert tls.Certificate
|
||||
certLoaded bool
|
||||
|
||||
Reference in New Issue
Block a user