only notify watched resource (#7592)

This commit is contained in:
Zexi Li
2020-08-19 12:52:26 +08:00
committed by GitHub
parent cc8c7b64a7
commit 86253534d7
10 changed files with 419 additions and 94 deletions
+21 -6
View File
@@ -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)
}
}
}
}
}
@@ -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)
}
@@ -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)
+103 -47
View File
@@ -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)
}
}
+190
View File
@@ -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
}
}
+46
View File
@@ -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
}
+13 -11
View File
@@ -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)
}