From 242d7eaa0f5cae16218ead5b49d1cbeb31f61030 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Wed, 25 Aug 2021 19:44:15 +0800 Subject: [PATCH] fix(cloudcommon): avoid service exit when etcd unreachable --- pkg/cloudcommon/database.go | 54 +++++++++++++++------- pkg/cloudcommon/syncman/watcher/watcher.go | 13 +++--- pkg/mcclient/informer/watcher.go | 21 +++++++++ 3 files changed, 65 insertions(+), 23 deletions(-) diff --git a/pkg/cloudcommon/database.go b/pkg/cloudcommon/database.go index bb1b307285..9cee1045a2 100644 --- a/pkg/cloudcommon/database.go +++ b/pkg/cloudcommon/database.go @@ -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() { diff --git a/pkg/cloudcommon/syncman/watcher/watcher.go b/pkg/cloudcommon/syncman/watcher/watcher.go index 39e3c47eb9..7f91419733 100644 --- a/pkg/cloudcommon/syncman/watcher/watcher.go +++ b/pkg/cloudcommon/syncman/watcher/watcher.go @@ -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 } diff --git a/pkg/mcclient/informer/watcher.go b/pkg/mcclient/informer/watcher.go index 4430b3b8ff..41dbc18b8a 100644 --- a/pkg/mcclient/informer/watcher.go +++ b/pkg/mcclient/informer/watcher.go @@ -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 {