mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #12016 from zexi/automated-cherry-pick-of-#12014-upstream-release-3.7
Automated cherry pick of #12014: fix(cloudcommon): avoid service exit when etcd unreachable
This commit is contained in:
+38
-16
@@ -19,9 +19,11 @@ import (
|
||||
"database/sql"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/sqlchemy"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/appsrv"
|
||||
@@ -86,25 +88,45 @@ func InitDB(options *common_options.DBOptions) {
|
||||
}
|
||||
// lm := lockman.NewNoopLockManager()
|
||||
|
||||
if len(options.EtcdEndpoints) != 0 {
|
||||
log.Infof("using etcd as resource informer backend")
|
||||
tlsCfg, err := options.GetEtcdTLSConfig()
|
||||
if err != nil {
|
||||
log.Fatalf("get etcd informer backend tls config err: %v", err)
|
||||
startInitInformer(options)
|
||||
}
|
||||
|
||||
// startInitInformer starts goroutine init informer backend
|
||||
func startInitInformer(options *common_options.DBOptions) {
|
||||
go func() {
|
||||
if len(options.EtcdEndpoints) == 0 {
|
||||
return
|
||||
}
|
||||
informerBackend, err := informer.NewEtcdBackend(&etcd.SEtcdOptions{
|
||||
EtcdEndpoint: options.EtcdEndpoints,
|
||||
EtcdTimeoutSeconds: 5,
|
||||
EtcdRequestTimeoutSeconds: 2,
|
||||
EtcdLeaseExpireSeconds: 5,
|
||||
EtcdEnabldSsl: options.EtcdUseTLS,
|
||||
TLSConfig: tlsCfg,
|
||||
}, nil)
|
||||
if err != nil {
|
||||
log.Fatalf("new etcd informer backend error: %v", err)
|
||||
for {
|
||||
log.Infof("using etcd as resource informer backend")
|
||||
if err := initInformer(options); err != nil {
|
||||
log.Errorf("Init informer error: %v", err)
|
||||
time.Sleep(10 * time.Second)
|
||||
} else {
|
||||
break
|
||||
}
|
||||
}
|
||||
informer.Init(informerBackend)
|
||||
}()
|
||||
}
|
||||
|
||||
func initInformer(options *common_options.DBOptions) error {
|
||||
tlsCfg, err := options.GetEtcdTLSConfig()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "get etcd informer backend tls config")
|
||||
}
|
||||
informerBackend, err := informer.NewEtcdBackend(&etcd.SEtcdOptions{
|
||||
EtcdEndpoint: options.EtcdEndpoints,
|
||||
EtcdTimeoutSeconds: 5,
|
||||
EtcdRequestTimeoutSeconds: 2,
|
||||
EtcdLeaseExpireSeconds: 5,
|
||||
EtcdEnabldSsl: options.EtcdUseTLS,
|
||||
TLSConfig: tlsCfg,
|
||||
}, nil)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "new etcd informer backend")
|
||||
}
|
||||
informer.Init(informerBackend)
|
||||
return nil
|
||||
}
|
||||
|
||||
func CloseDB() {
|
||||
|
||||
@@ -82,12 +82,11 @@ func (manager *SInformerSyncManager) startWatcher() error {
|
||||
log.Infof("[%s] Start resource informer watcher for %s", manager.Name(), manager.resourceManager.GetKeyword())
|
||||
ctx := context.Background()
|
||||
s := auth.GetAdminSession(ctx, consts.GetRegion(), "v1")
|
||||
watchMan, err := informer.NewWatchManagerBySession(s)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "NewWatchManagerBySession")
|
||||
}
|
||||
if err := watchMan.For(manager.resourceManager).AddEventHandler(ctx, manager); err != nil {
|
||||
return errors.Wrapf(err, "watch resource %s", manager.resourceManager.GetKeyword())
|
||||
}
|
||||
informer.NewWatchManagerBySessionBg(s, func(watchMan *informer.SWatchManager) error {
|
||||
if err := watchMan.For(manager.resourceManager).AddEventHandler(ctx, manager); err != nil {
|
||||
return errors.Wrapf(err, "watch resource %s", manager.resourceManager.GetKeyword())
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -16,8 +16,10 @@ package informer
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/errors"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/etcd"
|
||||
@@ -36,6 +38,25 @@ func NewWatchManagerBySession(session *mcclient.ClientSession) (*SWatchManager,
|
||||
return NewWatchManager(session.GetClient(), session.GetToken(), session.GetRegion(), session.GetEndpointType())
|
||||
}
|
||||
|
||||
func NewWatchManagerBySessionBg(session *mcclient.ClientSession, callback func(man *SWatchManager) error) {
|
||||
go func() {
|
||||
for {
|
||||
watchMan, err := NewWatchManagerBySession(session)
|
||||
if err != nil {
|
||||
log.Errorf("NewWatchManagerBySession error: %v", err)
|
||||
} else {
|
||||
if err := callback(watchMan); err != nil {
|
||||
log.Warningf("callback with watchMan error: %v", err)
|
||||
} else {
|
||||
log.Infof("callback with watchMan success.")
|
||||
break
|
||||
}
|
||||
}
|
||||
time.Sleep(10 * time.Second)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func NewWatchManager(client *mcclient.Client, token mcclient.TokenCredential, region, interfaceType string) (*SWatchManager, error) {
|
||||
endpoint, err := client.GetCommonEtcdEndpoint(token, region, interfaceType)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user