Replace AddEventListener with AddEventListenerCtx and Publish with PublishCtx (#42284)

This commit is contained in:
idafurjes
2021-11-29 14:23:24 +01:00
committed by GitHub
parent 8927a3ca20
commit a65e0be110
6 changed files with 17 additions and 62 deletions
+4 -49
View File
@@ -29,7 +29,7 @@ type TransactionManager interface {
type Bus interface {
Dispatch(msg Msg) error
DispatchCtx(ctx context.Context, msg Msg) error
Publish(msg Msg) error
PublishCtx(ctx context.Context, msg Msg) error
// InTransaction starts a transaction and store it in the context.
@@ -40,7 +40,7 @@ type Bus interface {
AddHandler(handler HandlerFunc)
AddHandlerCtx(handler HandlerFunc)
AddEventListener(handler HandlerFunc)
AddEventListenerCtx(handler HandlerFunc)
// SetTransactionManager allows the user to replace the internal
@@ -190,31 +190,6 @@ func (b *InProcBus) PublishCtx(ctx context.Context, msg Msg) error {
return nil
}
// Publish function publish a message to the bus listener.
func (b *InProcBus) Publish(msg Msg) error {
var msgName = reflect.TypeOf(msg).Elem().Name()
var params = []reflect.Value{}
if listeners, exists := b.listenersWithCtx[msgName]; exists {
params = append(params, reflect.ValueOf(context.Background()))
params = append(params, reflect.ValueOf(msg))
if setting.Env == setting.Dev {
b.logger.Warn("Publish called with message handler registered using AddEventHandlerCtx and should be changed to use PublishCtx", "msgName", msgName)
}
if err := callListeners(listeners, params); err != nil {
return err
}
}
if listeners, exists := b.listeners[msgName]; exists {
params = append(params, reflect.ValueOf(msg))
if err := callListeners(listeners, params); err != nil {
return err
}
}
return nil
}
func callListeners(listeners []HandlerFunc, params []reflect.Value) error {
for _, listenerHandler := range listeners {
ret := reflect.ValueOf(listenerHandler).Call(params)
@@ -247,16 +222,6 @@ func (b *InProcBus) GetHandlerCtx(name string) HandlerFunc {
return b.handlersWithCtx[name]
}
func (b *InProcBus) AddEventListener(handler HandlerFunc) {
handlerType := reflect.TypeOf(handler)
eventName := handlerType.In(0).Elem().Name()
_, exists := b.listeners[eventName]
if !exists {
b.listeners[eventName] = make([]HandlerFunc, 0)
}
b.listeners[eventName] = append(b.listeners[eventName], handler)
}
func (b *InProcBus) AddEventListenerCtx(handler HandlerFunc) {
handlerType := reflect.TypeOf(handler)
eventName := handlerType.In(1).Elem().Name()
@@ -279,12 +244,6 @@ func AddHandlerCtx(implName string, handler HandlerFunc) {
globalBus.AddHandlerCtx(handler)
}
// AddEventListener attaches a handler function to the event listener.
// Package level function.
func AddEventListener(handler HandlerFunc) {
globalBus.AddEventListener(handler)
}
// AddEventListenerCtx attaches a handler function to the event listener.
// Package level function.
func AddEventListenerCtx(handler HandlerFunc) {
@@ -299,12 +258,8 @@ func DispatchCtx(ctx context.Context, msg Msg) error {
return globalBus.DispatchCtx(ctx, msg)
}
func Publish(msg Msg) error {
return globalBus.Publish(msg)
}
func PublishCtx(msg Msg) error {
return globalBus.Publish(msg)
func PublishCtx(ctx context.Context, msg Msg) error {
return globalBus.PublishCtx(ctx, msg)
}
func GetHandlerCtx(name string) HandlerFunc {
+5 -5
View File
@@ -127,12 +127,12 @@ func TestEventPublish(t *testing.T) {
var invoked bool
bus.AddEventListener(func(query *testQuery) error {
bus.AddEventListenerCtx(func(ctx context.Context, query *testQuery) error {
invoked = true
return nil
})
err := bus.Publish(&testQuery{})
err := bus.PublishCtx(context.Background(), &testQuery{})
require.NoError(t, err, "unable to publish event")
require.True(t, invoked)
@@ -141,7 +141,7 @@ func TestEventPublish(t *testing.T) {
func TestEventPublish_NoRegisteredListener(t *testing.T) {
bus := New()
err := bus.Publish(&testQuery{})
err := bus.PublishCtx(context.Background(), &testQuery{})
require.NoError(t, err, "unable to publish event")
}
@@ -173,7 +173,7 @@ func TestEventPublishCtx(t *testing.T) {
var invoked bool
bus.AddEventListener(func(query *testQuery) error {
bus.AddEventListenerCtx(func(ctx context.Context, query *testQuery) error {
invoked = true
return nil
})
@@ -194,7 +194,7 @@ func TestEventCtxPublish(t *testing.T) {
return nil
})
err := bus.Publish(&testQuery{})
err := bus.PublishCtx(context.Background(), &testQuery{})
require.NoError(t, err, "unable to publish event")
require.True(t, invoked)