From ff53671a18bdd2cb175c9b6cdb24d313dd3b6293 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Thu, 29 Oct 2020 16:36:17 +0800 Subject: [PATCH] informer: ignore outside context cancel when worker run --- pkg/cloudcommon/informer/etcd.go | 4 ++++ pkg/cloudcommon/informer/etcd_client.go | 2 +- pkg/cloudcommon/informer/informer.go | 6 +++--- pkg/cloudcommon/informer/worker.go | 27 +++++++++++++++++++++++-- 4 files changed, 33 insertions(+), 6 deletions(-) diff --git a/pkg/cloudcommon/informer/etcd.go b/pkg/cloudcommon/informer/etcd.go index 3967a7122d..2c9e0da760 100644 --- a/pkg/cloudcommon/informer/etcd.go +++ b/pkg/cloudcommon/informer/etcd.go @@ -156,6 +156,10 @@ func (b *EtcdBackend) put(ctx context.Context, key, val string) error { return b.PutWithLease(ctx, key, val, b.leaseTTL) } +func (b *EtcdBackend) PutSession(ctx context.Context, key, val string) error { + return b.client.PutSession(ctx, key, val) +} + func (b *EtcdBackend) PutWithLease(ctx context.Context, key, val string, ttlSeconds int64) error { return b.client.PutWithLease(ctx, key, val, ttlSeconds) } diff --git a/pkg/cloudcommon/informer/etcd_client.go b/pkg/cloudcommon/informer/etcd_client.go index e03831b07c..72713c20c9 100644 --- a/pkg/cloudcommon/informer/etcd_client.go +++ b/pkg/cloudcommon/informer/etcd_client.go @@ -99,7 +99,7 @@ func (b *EtcdBackendForClient) registerClientResource(ctx context.Context, key s if err != nil { return err } - if err := b.PutWithLease(ctx, clientKey, "ok", 60); err != nil { + if err := b.PutSession(ctx, clientKey, "ok"); err != nil { return err } return nil diff --git a/pkg/cloudcommon/informer/informer.go b/pkg/cloudcommon/informer/informer.go index e338bcca4d..e56fe43587 100644 --- a/pkg/cloudcommon/informer/informer.go +++ b/pkg/cloudcommon/informer/informer.go @@ -92,7 +92,7 @@ func Create(ctx context.Context, obj *ModelObject) error { if !isResourceWatched(obj.KeywordPlural) { return nil } - return run(func(be IInformerBackend) error { + return run(ctx, func(ctx context.Context, be IInformerBackend) error { return be.Create(ctx, obj) }) } @@ -101,7 +101,7 @@ func Update(ctx context.Context, obj *ModelObject, oldObj *jsonutils.JSONDict) e if !isResourceWatched(obj.KeywordPlural) { return nil } - return run(func(be IInformerBackend) error { + return run(ctx, func(ctx context.Context, be IInformerBackend) error { return be.Update(ctx, obj, oldObj) }) } @@ -110,7 +110,7 @@ func Delete(ctx context.Context, obj *ModelObject) error { if !isResourceWatched(obj.KeywordPlural) { return nil } - return run(func(be IInformerBackend) error { + return run(ctx, func(ctx context.Context, be IInformerBackend) error { return be.Delete(ctx, obj) }) } diff --git a/pkg/cloudcommon/informer/worker.go b/pkg/cloudcommon/informer/worker.go index e85bcf5018..68afba9f94 100644 --- a/pkg/cloudcommon/informer/worker.go +++ b/pkg/cloudcommon/informer/worker.go @@ -15,6 +15,8 @@ package informer import ( + "context" + "yunion.io/x/log" "yunion.io/x/onecloud/pkg/appsrv" @@ -29,14 +31,35 @@ func init() { informerWorkerMan = appsrv.NewWorkerManager("InformerWorkerManager", 10, 10240, false) } -func run(f func(be IInformerBackend) error) error { +/*type noCancel struct { + ctx context.Context +} + +func (c noCancel) Deadline() (time.Time, bool) { + return time.Time{}, false +} + +func (c noCancel) Done() <-chan struct{} { + return nil +} + +func (c noCancel) Err() error { + return nil +} + +func (c noCancel) Value(key interface{}) interface{} { + return c.ctx.Value(key) +}*/ + +func run(ctx context.Context, f func(ctx context.Context, be IInformerBackend) error) error { be := GetDefaultBackend() if be == nil { return ErrBackendNotInit } wf := func() { nopanic.Run(func() { - if err := f(be); err != nil { + // outside context ignored cause of run in worker + if err := f(context.Background(), be); err != nil { log.Errorf("run informer error: %v", err) } })