diff --git a/cmd/climc/shell/watch.go b/cmd/climc/shell/watch.go index 9ad5eefbcc..a68e4e9793 100644 --- a/cmd/climc/shell/watch.go +++ b/cmd/climc/shell/watch.go @@ -54,7 +54,7 @@ func init() { } R(&WatchOptions{}, "watch", "Watch resources", func(s *mcclient.ClientSession, opts *WatchOptions) error { - watchMan, err := informer.NewWatchManagerBySession(s, nil) + watchMan, err := informer.NewWatchManagerBySession(s) if err != nil { return err } diff --git a/pkg/cloudcommon/etcd/etcd.go b/pkg/cloudcommon/etcd/etcd.go index 74ce98a06e..ced633d184 100644 --- a/pkg/cloudcommon/etcd/etcd.go +++ b/pkg/cloudcommon/etcd/etcd.go @@ -21,6 +21,7 @@ import ( "time" "go.etcd.io/etcd/clientv3" + "go.etcd.io/etcd/mvcc/mvccpb" "google.golang.org/grpc" "yunion.io/x/log" @@ -161,7 +162,7 @@ func (cli *SEtcdClient) startSession() error { for { if _, ok := <-ch; !ok { cli.leaseLiving = false - log.Errorf("fail to keepalive sessoin") + log.Errorf("fail to keepalive session") if cli.onKeepaliveFailure != nil { cli.onKeepaliveFailure() } @@ -281,8 +282,9 @@ func (cli *SEtcdClient) List(ctx context.Context, prefix string) ([]SEtcdKeyValu return ret, nil } -type TEtcdCreateEventFunc func(key, value []byte) -type TEtcdModifyEventFunc func(key, oldvalue, value []byte) +type TEtcdCreateEventFunc func(ctx context.Context, key, value []byte) +type TEtcdModifyEventFunc func(ctx context.Context, key, oldvalue, value []byte) +type TEtcdDeleteEventFunc func(ctx context.Context, key []byte) type SEtcdWatcher struct { watcher clientv3.Watcher @@ -294,7 +296,12 @@ func (w *SEtcdWatcher) Cancel() { w.cancel() } -func (cli *SEtcdClient) Watch(ctx context.Context, prefix string, onCreate TEtcdCreateEventFunc, onModify TEtcdModifyEventFunc) error { +func (cli *SEtcdClient) Watch( + ctx context.Context, prefix string, + onCreate TEtcdCreateEventFunc, + onModify TEtcdModifyEventFunc, + onDelete TEtcdDeleteEventFunc, +) error { _, ok := cli.watchers[prefix] if ok { return errors.Errorf("watch prefix %s already registered", prefix) @@ -314,10 +321,18 @@ func (cli *SEtcdClient) Watch(ctx context.Context, prefix string, onCreate TEtcd go func() { for wresp := range rch { for _, ev := range wresp.Events { + key := ev.Kv.Key[len(cli.namespace):] if ev.PrevKv == nil { - onCreate(ev.Kv.Key[len(cli.namespace):], ev.Kv.Value) + onCreate(nctx, key, ev.Kv.Value) } else { - onModify(ev.Kv.Key[len(cli.namespace):], ev.PrevKv.Value, ev.Kv.Value) + switch ev.Type { + case mvccpb.PUT: + onModify(nctx, key, ev.PrevKv.Value, ev.Kv.Value) + case mvccpb.DELETE: + if onDelete != nil { + onDelete(nctx, key) + } + } } } } diff --git a/pkg/cloudcommon/etcd/models/base/basemanager.go b/pkg/cloudcommon/etcd/models/base/basemanager.go index 9f6a4c2467..68451e9386 100644 --- a/pkg/cloudcommon/etcd/models/base/basemanager.go +++ b/pkg/cloudcommon/etcd/models/base/basemanager.go @@ -234,7 +234,8 @@ func (manager *SEtcdBaseModelManager) Delete(ctx context.Context, model IEtcdMod func (manager *SEtcdBaseModelManager) Watch(ctx context.Context, onCreate etcd.TEtcdCreateEventFunc, onModify etcd.TEtcdModifyEventFunc, + onDelete etcd.TEtcdDeleteEventFunc, ) { prefix := manager.managerKey() - etcd.Default().Watch(ctx, prefix, onCreate, onModify) + etcd.Default().Watch(ctx, prefix, onCreate, onModify, onDelete) } diff --git a/pkg/cloudcommon/etcd/models/base/interface.go b/pkg/cloudcommon/etcd/models/base/interface.go index 800097dae2..7e6420bfcf 100644 --- a/pkg/cloudcommon/etcd/models/base/interface.go +++ b/pkg/cloudcommon/etcd/models/base/interface.go @@ -40,7 +40,7 @@ type IEtcdModelManager interface { Save(ctx context.Context, model IEtcdModel) error Delete(ctx context.Context, model IEtcdModel) error Session(ctx context.Context, model IEtcdModel) error - Watch(ctx context.Context, onCreate etcd.TEtcdCreateEventFunc, onModify etcd.TEtcdModifyEventFunc) + Watch(ctx context.Context, onCreate etcd.TEtcdCreateEventFunc, onModify etcd.TEtcdModifyEventFunc, onDelete etcd.TEtcdDeleteEventFunc) CustomizeHandlerInfo(handler *appsrv.SHandlerInfo) FetchCreateHeaderData(ctx context.Context, header http.Header) (jsonutils.JSONObject, error) diff --git a/pkg/cloudcommon/informer/etcd.go b/pkg/cloudcommon/informer/etcd.go index 6d43d60f40..497ad1d9d9 100644 --- a/pkg/cloudcommon/informer/etcd.go +++ b/pkg/cloudcommon/informer/etcd.go @@ -18,6 +18,7 @@ import ( "context" "fmt" "path/filepath" + "strings" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -27,7 +28,8 @@ import ( ) const ( - EtcdInformerPrefix = "/onecloud/informer" + EtcdInformerPrefix = "/onecloud/informer" + EtcdInformerClientsKey = "@clients" EventTypeCreate = "CREATE" EventTypeUpdate = "UPDATE" @@ -63,7 +65,7 @@ type EtcdBackend struct { leaseTTL int64 } -func NewEtcdBackend(opt *etcd.SEtcdOptions, onKeepaliveFailure func()) (*EtcdBackend, error) { +func newEtcdBackend(opt *etcd.SEtcdOptions, onKeepaliveFailure func()) (*EtcdBackend, error) { opt.EtcdNamspace = EtcdInformerPrefix be := new(EtcdBackend) be.leaseTTL = int64(opt.EtcdLeaseExpireSeconds) @@ -78,6 +80,31 @@ func NewEtcdBackend(opt *etcd.SEtcdOptions, onKeepaliveFailure func()) (*EtcdBac return be, nil } +func NewEtcdBackend(opt *etcd.SEtcdOptions, onKeepaliveFailure func()) (*EtcdBackend, error) { + be, err := newEtcdBackend(opt, onKeepaliveFailure) + if err != nil { + return nil, err + } + ctx := context.Background() + be.initClientResources(ctx) + be.StartClientWatch(ctx) + return be, nil +} + +func (b *EtcdBackend) initClientResources(ctx context.Context) error { + pairs, err := b.client.List(ctx, "/") + if err != nil { + return errors.Wrap(err, "list all client resources") + } + for _, pair := range pairs { + keywordPlural, err := b.getClientWatchResource([]byte(pair.Key)) + if err == nil && len(keywordPlural) != 0 { + AddWatchedResources(keywordPlural) + } + } + return nil +} + func (b *EtcdBackend) getObjectKey(obj *ModelObject) string { if obj.IsJoint { return fmt.Sprintf("/%s/%s/%s", obj.KeywordPlural, obj.MasterId, obj.SlaveId) @@ -126,59 +153,88 @@ func (b *EtcdBackend) Delete(ctx context.Context, obj *ModelObject) error { } func (b *EtcdBackend) put(ctx context.Context, key, val string) error { - return b.client.PutWithLease(ctx, key, val, b.leaseTTL) + return b.PutWithLease(ctx, key, val, b.leaseTTL) +} + +func (b *EtcdBackend) PutWithLease(ctx context.Context, key, val string, ttlSeconds int64) error { + return b.client.PutWithLease(ctx, key, val, ttlSeconds) } func (b *EtcdBackend) onKeepaliveFailure() { if err := b.client.RestartSession(); err != nil { log.Errorf("restart etcd session error: %v", err) + return + } + b.StartClientWatch(context.Background()) +} + +func (b *EtcdBackend) StartClientWatch(ctx context.Context) { + b.client.Unwatch("/") + b.client.Watch(ctx, "/", b.onClientResourceCreate, b.onClientResourceUpdate, b.onClientResourceDelete) +} + +func (b *EtcdBackend) isClientsKey(key []byte) bool { + return strings.Contains(string(key), EtcdInformerClientsKey) +} + +func (b *EtcdBackend) getClientWatchResource(key []byte) (string, error) { + // key is like: /servers/@clients/default-climc-5d4c8d49f6-p6l68 + keyPath := string(key) + parts := strings.Split(keyPath, "/") + if len(parts) != 4 { + return "", errors.Errorf("invalid client resource key: %v", parts) + } + if parts[2] != EtcdInformerClientsKey { + return "", errors.Errorf("key %s not contains %s", keyPath, EtcdInformerClientsKey) + } + return parts[1], nil +} + +func (b *EtcdBackend) onClientResourceAdd(key []byte) { + if !b.isClientsKey(key) { + return + } + keywordPlural, err := b.getClientWatchResource(key) + if err != nil { + log.Errorf("get client watch resource error: %v", err) + return + } + AddWatchedResources(keywordPlural) +} + +func (b *EtcdBackend) onClientResourceCreate(ctx context.Context, key, value []byte) { + b.onClientResourceAdd(key) +} + +func (b *EtcdBackend) onClientResourceUpdate(ctx context.Context, key, oldvalue, value []byte) { + b.onClientResourceAdd(key) +} + +func (b *EtcdBackend) shouldDeleteWatchedResource(ctx context.Context, keywordPlural string) bool { + // clientsKey is like: /servers/@clients + clientsKey := fmt.Sprintf("/%s/%s", keywordPlural, EtcdInformerClientsKey) + pairs, err := b.client.List(ctx, clientsKey) + if err != nil { + log.Errorf("list clientsKey %s error: %v", err) + return false + } + return len(pairs) == 0 +} + +func (b *EtcdBackend) onClientResourceDelete(ctx context.Context, key []byte) { + if !b.isClientsKey(key) { + return + } + keywordPlural, err := b.getClientWatchResource(key) + if err != nil { + log.Errorf("get client watch resource keywordPlural error: %v", err) + return + } + if b.shouldDeleteWatchedResource(ctx, keywordPlural) { + DeleteWatchedResources(keywordPlural) } } func (b *EtcdBackend) getWatchKey(key string) string { return filepath.Join("/", key) } - -func (b *EtcdBackend) Watch(ctx context.Context, key string, handler ResourceEventHandler) error { - return b.client.Watch(ctx, b.getWatchKey(key), b.onCreate(handler), b.onModify(handler)) -} - -func (b *EtcdBackend) Unwatch(key string) { - b.client.Unwatch(b.getWatchKey(key)) -} - -func (b *EtcdBackend) onCreate(handler ResourceEventHandler) etcd.TEtcdCreateEventFunc { - return func(key, value []byte) { - b.processEvent(handler, string(key), value) - } -} - -func (b *EtcdBackend) onModify(handler ResourceEventHandler) etcd.TEtcdModifyEventFunc { - return func(key, _, value []byte) { - // not care about oldvalue, so ignore it - b.processEvent(handler, string(key), value) - } -} - -func (b *EtcdBackend) processEvent(handler ResourceEventHandler, key string, value []byte) { - if len(value) == 0 { - // object already deleted by lease out of ttl - return - } - mObj, err := newModelObjectFromValue(value) - if err != nil { - log.Errorf("new %s model objecd from value error: %v", key, err) - return - } - eType := mObj.EventType - switch eType { - case EventTypeCreate: - handler.OnAdd(mObj.Object) - case EventTypeUpdate: - handler.OnUpdate(mObj.OldObject, mObj.Object) - case EventTypeDelete: - handler.OnDelete(mObj.Object) - default: - log.Errorf("Invalid event type: %s, mObj: %#v", eType, mObj) - } -} diff --git a/pkg/cloudcommon/informer/etcd_client.go b/pkg/cloudcommon/informer/etcd_client.go new file mode 100644 index 0000000000..e03831b07c --- /dev/null +++ b/pkg/cloudcommon/informer/etcd_client.go @@ -0,0 +1,190 @@ +// 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 informer + +import ( + "context" + "fmt" + "os" + "path/filepath" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/cloudcommon/etcd" +) + +type IWatcher interface { + Watch(ctx context.Context, key string, handler ResourceEventHandler) error + Unwatch(key string) +} + +type ResourceEventHandler interface { + OnAdd(obj *jsonutils.JSONDict) + OnUpdate(oldObj, newObj *jsonutils.JSONDict) + OnDelete(obj *jsonutils.JSONDict) +} + +type EtcdBackendForClient struct { + *EtcdBackend + clientResources map[string]ResourceEventHandler +} + +func NewEtcdBackendForClient(opt *etcd.SEtcdOptions) (*EtcdBackendForClient, error) { + bec := &EtcdBackendForClient{ + clientResources: make(map[string]ResourceEventHandler, 0), + } + be, err := newEtcdBackend(opt, bec.onKeepaliveFailure) + if err != nil { + return nil, err + } + bec.EtcdBackend = be + bec.StartClientWatch(context.Background()) + return bec, nil +} + +func (b *EtcdBackendForClient) onKeepaliveFailure() { + if err := b.client.RestartSession(); err != nil { + log.Errorf("restart etcd session error: %v", err) + return + } + b.StartClientWatch(context.Background()) +} + +func (b *EtcdBackendForClient) StartClientWatch(ctx context.Context) { + wf := func(key string) string { + return filepath.Join(EtcdInformerPrefix, key) + } + for key, handler := range b.clientResources { + // if etcd pod deleted, should unwatch then rewatch resource + b.Unwatch(key) + if err := b.Watch(ctx, key, handler); err != nil { + log.Errorf("start watch client resource %s error: %v", key, err) + continue + } + log.Infof("%s rewatched", wf(key)) + } + b.client.Unwatch("/") + if err := b.client.Watch(ctx, "/", b.onClientResourceCreate, b.onClientResourceUpdate, b.onClientResourceDelete); err != nil { + log.Errorf("start watch %s error: %v", wf("/"), err) + } else { + log.Infof("%s watched", wf("/")) + } +} + +func (b *EtcdBackend) getClientRegisterKey(resKey string) (string, error) { + hostname, err := os.Hostname() + if err != nil { + return "", errors.Wrap(err, "get hostname") + } + clientKey := fmt.Sprintf("/%s/%s/%s", resKey, EtcdInformerClientsKey, hostname) + return clientKey, nil +} + +func (b *EtcdBackendForClient) registerClientResource(ctx context.Context, key string) error { + clientKey, err := b.getClientRegisterKey(key) + if err != nil { + return err + } + if err := b.PutWithLease(ctx, clientKey, "ok", 60); err != nil { + return err + } + return nil +} + +func (b *EtcdBackendForClient) Watch(ctx context.Context, key string, handler ResourceEventHandler) error { + b.clientResources[key] = handler + if err := b.registerClientResource(ctx, key); err != nil { + return errors.Wrapf(err, "register watch client resource %s", key) + } + return b.client.Watch(ctx, b.getWatchKey(key), b.onCreate(ctx, handler), b.onModify(ctx, handler), nil) +} + +func (b *EtcdBackendForClient) Unwatch(key string) { + delete(b.clientResources, key) + b.client.Unwatch(b.getWatchKey(key)) +} + +func (b *EtcdBackendForClient) onCreate(ctx context.Context, handler ResourceEventHandler) etcd.TEtcdCreateEventFunc { + return func(ctx context.Context, key, value []byte) { + b.processEvent(handler, string(key), value) + } +} + +func (b *EtcdBackendForClient) onModify(ctx context.Context, handler ResourceEventHandler) etcd.TEtcdModifyEventFunc { + return func(ctx context.Context, key, _, value []byte) { + // not care about oldvalue, so ignore it + b.processEvent(handler, string(key), value) + } +} + +func (b *EtcdBackendForClient) processEvent(handler ResourceEventHandler, key string, value []byte) { + if len(value) == 0 { + // object already deleted by lease out of ttl + return + } + if b.isClientsKey([]byte(key)) { + return + } + mObj, err := newModelObjectFromValue(value) + if err != nil { + log.Errorf("new %s model objecd from value error: %v", key, err) + return + } + eType := mObj.EventType + switch eType { + case EventTypeCreate: + handler.OnAdd(mObj.Object) + case EventTypeUpdate: + handler.OnUpdate(mObj.OldObject, mObj.Object) + case EventTypeDelete: + handler.OnDelete(mObj.Object) + default: + log.Errorf("Invalid key %s, event type: %s, mObj: %#v", key, eType, mObj) + } +} + +func (b *EtcdBackendForClient) onClientResourceAdd(key []byte) { + // do nothing + // self registered key has store in b.clientResources +} + +func (b *EtcdBackendForClient) onClientResourceCreate(ctx context.Context, key, value []byte) { + b.onClientResourceAdd(key) +} + +func (b *EtcdBackendForClient) onClientResourceUpdate(ctx context.Context, key, oldvalue, value []byte) { + b.onClientResourceAdd(key) +} + +func (b *EtcdBackendForClient) onClientResourceDelete(ctx context.Context, key []byte) { + if !b.isClientsKey(key) { + return + } + keywordPlural, err := b.getClientWatchResource(key) + if err != nil { + log.Errorf("get client watch resource keywordPlural error: %v", err) + return + } + _, ok := b.clientResources[keywordPlural] + if !ok { + return + } + if err := b.registerClientResource(ctx, keywordPlural); err != nil { + log.Errorf("rewatch client resource %s error: %v", keywordPlural, err) + return + } +} diff --git a/pkg/cloudcommon/informer/handler.go b/pkg/cloudcommon/informer/handler.go new file mode 100644 index 0000000000..dd4f1697c6 --- /dev/null +++ b/pkg/cloudcommon/informer/handler.go @@ -0,0 +1,46 @@ +// 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 informer + +import ( + "sync" + + "yunion.io/x/pkg/util/sets" +) + +var ( + globalWatchResources = new(sync.Map) +) + +func AddWatchedResources(resources ...string) { + for _, res := range resources { + globalWatchResources.Store(res, true) + } +} + +func DeleteWatchedResources(resources ...string) { + for _, res := range resources { + globalWatchResources.Delete(res) + } +} + +func GetWatchResources() sets.String { + ret := sets.NewString() + globalWatchResources.Range(func(key, val interface{}) bool { + ret.Insert(key.(string)) + return true + }) + return ret +} diff --git a/pkg/cloudcommon/informer/informer.go b/pkg/cloudcommon/informer/informer.go index ead1f0866d..e338bcca4d 100644 --- a/pkg/cloudcommon/informer/informer.go +++ b/pkg/cloudcommon/informer/informer.go @@ -37,11 +37,6 @@ type IInformerBackend interface { Delete(ctx context.Context, obj *ModelObject) error } -type IWatcher interface { - Watch(ctx context.Context, key string, handler ResourceEventHandler) error - Unwatch(key string) -} - func Init(be IInformerBackend) { if defaultBackend != nil { log.Fatalf("informer backend %q already init", be.GetType()) @@ -89,26 +84,33 @@ func NewJointModel(obj interface{}, keywordPlural, masterId, slaveId string) *Mo return model } +func isResourceWatched(keywordPlural string) bool { + return GetWatchResources().Has(keywordPlural) +} + func Create(ctx context.Context, obj *ModelObject) error { + if !isResourceWatched(obj.KeywordPlural) { + return nil + } return run(func(be IInformerBackend) error { return be.Create(ctx, obj) }) } func Update(ctx context.Context, obj *ModelObject, oldObj *jsonutils.JSONDict) error { + if !isResourceWatched(obj.KeywordPlural) { + return nil + } return run(func(be IInformerBackend) error { return be.Update(ctx, obj, oldObj) }) } func Delete(ctx context.Context, obj *ModelObject) error { + if !isResourceWatched(obj.KeywordPlural) { + return nil + } return run(func(be IInformerBackend) error { return be.Delete(ctx, obj) }) } - -type ResourceEventHandler interface { - OnAdd(obj *jsonutils.JSONDict) - OnUpdate(oldObj, newObj *jsonutils.JSONDict) - OnDelete(obj *jsonutils.JSONDict) -} diff --git a/pkg/compute/models/host_health.go b/pkg/compute/models/host_health.go index 69f95e2227..56a07e7148 100644 --- a/pkg/compute/models/host_health.go +++ b/pkg/compute/models/host_health.go @@ -100,7 +100,7 @@ func (h *SHostHealthChecker) startWatcher(ctx context.Context, hostId string) { if _, ok := h.hc[hostId]; !ok { h.hc[hostId] = ch } - h.cli.Watch(ctx, key, h.onHostOnline(hostId), h.onHostOffline(hostId)) + h.cli.Watch(ctx, key, h.onHostOnline(ctx, hostId), h.onHostOffline(ctx, hostId), h.onHostOfflineDeleted(ctx, hostId)) } func (h *SHostHealthChecker) onHostUnhealthy(ctx context.Context, hostId string) { @@ -112,8 +112,8 @@ func (h *SHostHealthChecker) onHostUnhealthy(ctx context.Context, hostId string) } } -func (h *SHostHealthChecker) onHostOnline(hostId string) etcd.TEtcdCreateEventFunc { - return func(key, value []byte) { +func (h *SHostHealthChecker) onHostOnline(ctx context.Context, hostId string) etcd.TEtcdCreateEventFunc { + return func(ctx context.Context, key, value []byte) { log.Infof("Got host online %s", hostId) if h.hc[hostId] != nil { h.hc[hostId] <- struct{}{} @@ -121,17 +121,27 @@ func (h *SHostHealthChecker) onHostOnline(hostId string) etcd.TEtcdCreateEventFu } } -func (h *SHostHealthChecker) onHostOffline(hostId string) etcd.TEtcdModifyEventFunc { - return func(key, oldvalue, value []byte) { - log.Warningf("host %s disconnect with etcd", hostId) - go func() { - select { - case <-time.NewTimer(h.timeout).C: - h.onHostUnhealthy(context.Background(), hostId) - case <-h.hc[hostId]: - h.startWatcher(context.Background(), hostId) - } - }() +func (h *SHostHealthChecker) processHostOffline(ctx context.Context, hostId string) { + log.Warningf("host %s disconnect with etcd", hostId) + go func() { + select { + case <-time.NewTimer(h.timeout).C: + h.onHostUnhealthy(ctx, hostId) + case <-h.hc[hostId]: + h.startWatcher(ctx, hostId) + } + }() +} + +func (h *SHostHealthChecker) onHostOffline(ctx context.Context, hostId string) etcd.TEtcdModifyEventFunc { + return func(ctx context.Context, key, oldvalue, value []byte) { + h.processHostOffline(ctx, hostId) + } +} + +func (h *SHostHealthChecker) onHostOfflineDeleted(ctx context.Context, hostId string) etcd.TEtcdDeleteEventFunc { + return func(ctx context.Context, key []byte) { + h.processHostOffline(ctx, hostId) } } diff --git a/pkg/mcclient/informer/watcher.go b/pkg/mcclient/informer/watcher.go index dd80db327c..4430b3b8ff 100644 --- a/pkg/mcclient/informer/watcher.go +++ b/pkg/mcclient/informer/watcher.go @@ -26,15 +26,17 @@ import ( ) type SWatchManager struct { - client *mcclient.Client - watchBackend informer.IWatcher + client *mcclient.Client + region string + interfaceType string + watchBackend informer.IWatcher } -func NewWatchManagerBySession(session *mcclient.ClientSession, onKeepaliveFailure func()) (*SWatchManager, error) { - return NewWatchManager(session.GetClient(), session.GetToken(), session.GetRegion(), session.GetEndpointType(), onKeepaliveFailure) +func NewWatchManagerBySession(session *mcclient.ClientSession) (*SWatchManager, error) { + return NewWatchManager(session.GetClient(), session.GetToken(), session.GetRegion(), session.GetEndpointType()) } -func NewWatchManager(client *mcclient.Client, token mcclient.TokenCredential, region, interfaceType string, onKeepaliveFailure func()) (*SWatchManager, error) { +func NewWatchManager(client *mcclient.Client, token mcclient.TokenCredential, region, interfaceType string) (*SWatchManager, error) { endpoint, err := client.GetCommonEtcdEndpoint(token, region, interfaceType) if err != nil { return nil, errors.Wrap(err, "get common etcd endpoint") @@ -54,13 +56,15 @@ func NewWatchManager(client *mcclient.Client, token mcclient.TokenCredential, re opt.EtcdEnabldSsl = true opt.TLSConfig = tlsCfg } - be, err := informer.NewEtcdBackend(opt, onKeepaliveFailure) + be, err := informer.NewEtcdBackendForClient(opt) if err != nil { return nil, errors.Wrap(err, "new etcd informer backend") } man := &SWatchManager{ - client: client, - watchBackend: be, + client: client, + region: region, + interfaceType: interfaceType, + watchBackend: be, } return man, nil } @@ -87,11 +91,12 @@ type sWatcher struct { eventHandler EventHandler } -func (man *SWatchManager) For(resourceManager IResourceManager) IWatcher { - return &sWatcher{ +func (man *SWatchManager) For(resMan IResourceManager) IWatcher { + watcher := &sWatcher{ manager: man, - resourceManager: resourceManager, + resourceManager: resMan, } + return watcher } func (w *sWatcher) AddEventHandler(ctx context.Context, handler EventHandler) error { @@ -100,8 +105,8 @@ func (w *sWatcher) AddEventHandler(ctx context.Context, handler EventHandler) er return w.manager.watch(w.ctx, w.resourceManager, w.eventHandler) } -func (man *SWatchManager) watch(ctx context.Context, resourceManager IResourceManager, handler informer.ResourceEventHandler) error { - return man.watchBackend.Watch(ctx, resourceManager.KeyString(), handler) +func (man *SWatchManager) watch(ctx context.Context, resMan IResourceManager, handler informer.ResourceEventHandler) error { + return man.watchBackend.Watch(ctx, resMan.KeyString(), handler) } func (w *sWatcher) wrapEventHandler(handler EventHandler) informer.ResourceEventHandler {