mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #8558 from zexi/hotfix/informer-context-cancelled
informer: ignore outside context cancel when worker run
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user