events API

events

package

API reference for the events package.

T
type

Handler

Handler processes an event of type T.

core/events/bus.go:16-16
type Handler func(ctx context.Context, event T) error
T
type

Priority

Priority defines the ordering of event handlers.

core/events/bus.go:19-19
type Priority int
T
type

DispatchStrategy

DispatchStrategy controls error handling during event dispatch.

core/events/bus.go:31-31
type DispatchStrategy int
T
type

Middleware

Middleware wraps event dispatch with cross-cutting behavior.

core/events/bus.go:41-41
type Middleware func(ctx context.Context, event any, next func(ctx context.Context, event any) error) error
S
struct

asyncEvent

core/events/bus.go:43-47
type asyncEvent struct

Fields

Name Type Description
event any
ctx context.Context
emit func(ctx context.Context, event any) error
T
type

OverflowStrategy

OverflowStrategy controls behavior when the async channel is full.

core/events/bus.go:50-50
type OverflowStrategy int
S
struct

Bus

Bus is the event bus that dispatches events to registered handlers.

core/events/bus.go:60-78
type Bus struct

Methods

func (*Bus) asyncProcessor()
{
	defer close(b.asyncDone)
	for {
		select {
		case <-b.asyncClose:
			for {
				select {
				case evt := <-b.asyncCh:
					b.processAsync(evt)
				default:
					return
				}
			}
		case evt := <-b.asyncCh:
			b.processAsync(evt)
		}
	}
}
processAsync
Method

Parameters

func (*Bus) processAsync(evt asyncEvent)
{
	err := func() (err error) {
		defer func() {
			if recovered := recover(); recovered != nil {
				err = fmt.Errorf("events: async handler panic: %v", recovered)
			}
		}()
		ctx := context.WithValue(evt.ctx, asyncBusContextKey{}, b)
		return evt.emit(ctx, evt.event)
	}()
	if err != nil {
		b.reportAsyncError(err)
	}
}

Parameters

err error
func (*Bus) reportAsyncError(err error)
{
	b.mu.RLock()
	fn := b.onAsyncError
	b.mu.RUnlock()
	if fn == nil {
		return
	}
	func() {
		defer func() {
			_ = recover()
		}()
		fn(err)
	}()
}
Close
Method

Close starts an asynchronous shutdown and never blocks the calling handler.

func (*Bus) Close()
{
	b.closeOnce.Do(func() {
		b.asyncMu.Lock()
		b.asyncClosed = true
		b.asyncMu.Unlock()
		go func() {
			b.asyncWG.Wait()
			close(b.asyncClose)
			<-b.asyncDone
			close(b.shutdownDone)
		}()
	})
}
Shutdown
Method

Shutdown drains queued events and waits for active handlers.

Parameters

Returns

error
func (*Bus) Shutdown(ctx context.Context) error
{
	if ctx == nil {
		ctx = context.Background()
	}
	if current, _ := ctx.Value(asyncBusContextKey{}).(*Bus); current == b {
		b.Close()
		return ErrReentrantShutdown
	}
	b.Close()
	select {
	case <-b.shutdownDone:
		return nil
	case <-ctx.Done():
		return ctx.Err()
	}
}
Use
Method

Use appends middleware to the bus.

Parameters

func (*Bus) Use(mw Middleware)
{
	b.mu.Lock()
	defer b.mu.Unlock()
	b.middlewares = append(b.middlewares, mw)
}

Fields

Name Type Description
subscribers *safemap.Map[reflect.Type, []subscriber]
strategy DispatchStrategy
middlewares []Middleware
onAsyncError func(error)
wildcard []subscriber
mu sync.RWMutex
asyncCh chan asyncEvent
asyncClose chan struct{}
asyncDone chan struct{}
asyncMu sync.Mutex
asyncWG sync.WaitGroup
asyncClosed bool
overflowStrat OverflowStrategy
bufferSize int
closeOnce sync.Once
shutdownDone chan struct{}
S
struct

subscriber

core/events/bus.go:80-83
type subscriber struct

Fields

Name Type Description
handler any
priority Priority
S
struct

asyncBusContextKey

core/events/bus.go:85-85
type asyncBusContextKey struct
F
function

Default

Default returns the package-level default Bus.

Returns

core/events/bus.go:99-101
func Default() *Bus

{
	return defaultBus
}
T
type

Option

Option configures a Bus.

core/events/bus.go:104-104
type Option options.Option[Bus]
F
function

New

New creates a new Bus with the given options.

Parameters

opts
...Option

Returns

core/events/bus.go:107-121
func New(opts ...Option) *Bus

{
	b := &Bus{
		subscribers:   safemap.New[reflect.Type, []subscriber](),
		strategy:      StopOnFirstError,
		asyncCh:       make(chan asyncEvent, 1024),
		asyncClose:    make(chan struct{}),
		asyncDone:     make(chan struct{}),
		shutdownDone:  make(chan struct{}),
		bufferSize:    1024,
		overflowStrat: OverflowFail,
	}
	options.Apply(b, opts...)
	go b.asyncProcessor()
	return b
}
F
function

WithBufferSize

WithBufferSize sets the async channel buffer size.

Parameters

n
int

Returns

core/events/bus.go:124-132
func WithBufferSize(n int) Option

{
	return func(b *Bus) {
		if n <= 0 {
			panic("events: buffer size must be positive")
		}
		b.bufferSize = n
		b.asyncCh = make(chan asyncEvent, n)
	}
}
F
function

WithOverflowStrategy

WithOverflowStrategy sets the overflow behavior for async events.

Parameters

Returns

core/events/bus.go:135-142
func WithOverflowStrategy(s OverflowStrategy) Option

{
	return func(b *Bus) {
		if s != OverflowFail && s != OverflowDropOldest {
			panic("events: invalid overflow strategy")
		}
		b.overflowStrat = s
	}
}
F
function

WithStrategy

WithStrategy sets the dispatch strategy for the Bus.

Parameters

Returns

core/events/bus.go:227-229
func WithStrategy(s DispatchStrategy) Option

{
	return func(b *Bus) { b.strategy = s }
}
F
function

WithOnAsyncError

WithOnAsyncError sets the error handler for async event emissions.

Parameters

fn
func(error)

Returns

core/events/bus.go:232-234
func WithOnAsyncError(fn func(error)) Option

{
	return func(b *Bus) { b.onAsyncError = fn }
}
F
function

Subscribe

Subscribe registers a typed handler for events of type T.

Parameters

b
fn
Handler[T]
priority
...Priority
core/events/bus.go:244-261
func Subscribe[T any](b *Bus, fn Handler[T], priority ...Priority)

{
	if b == nil {
		b = defaultBus
	}
	p := PriorityNormal
	if len(priority) > 0 {
		p = priority[0]
	}
	key := reflect.TypeFor[T]()
	b.subscribers.Compute(key, func(subs []subscriber, exists bool) []subscriber {
		newSubs := append([]subscriber(nil), subs...)
		newSubs = append(newSubs, subscriber{handler: fn, priority: p})
		sort.SliceStable(newSubs, func(i, j int) bool {
			return newSubs[i].priority > newSubs[j].priority
		})
		return newSubs
	})
}
F
function

SubscribeWildcard

SubscribeWildcard registers a handler for all event types.

Parameters

b
fn
func(ctx context.Context, event any) error
core/events/bus.go:264-271
func SubscribeWildcard(b *Bus, fn func(ctx context.Context, event any) error)

{
	if b == nil {
		b = defaultBus
	}
	b.mu.Lock()
	defer b.mu.Unlock()
	b.wildcard = append(b.wildcard, subscriber{handler: fn})
}
F
function

Emit

Emit dispatches an event to all matching handlers synchronously.

Parameters

b
event
T

Returns

error
core/events/bus.go:274-330
func Emit[T any](ctx context.Context, b *Bus, event T) error

{
	if b == nil {
		b = defaultBus
	}
	key := reflect.TypeFor[T]()
	subs, ok := b.subscribers.Get(key)

	b.mu.RLock()
	mws := b.middlewares
	b.mu.RUnlock()

	emit := func(ctx context.Context, evt any) error {
		if !ok {
			return nil
		}
		var errs []error
		for _, sub := range subs {
			if fn, ok := sub.handler.(Handler[T]); ok {
				if err := fn(ctx, evt.(T)); err != nil {
					if b.strategy == StopOnFirstError {
						return err
					}
					errs = append(errs, err)
				}
			}
		}
		if len(errs) > 0 {
			return errors.Join(errs...)
		}
		return nil
	}

	if len(mws) > 0 {
		chain := applyMiddleware(emit, mws)
		if err := chain(ctx, event); err != nil {
			return err
		}
		return emitWildcards(ctx, b, event)
	}

	if err := emit(ctx, event); err != nil {
		return err
	}

	b.mu.RLock()
	wildcards := b.wildcard
	b.mu.RUnlock()
	for _, w := range wildcards {
		if fn, ok := w.handler.(func(ctx context.Context, event any) error); ok {
			if err := fn(ctx, event); err != nil {
				return err
			}
		}
	}

	return nil
}
F
function

EmitAny

EmitAny dispatches an event using its runtime type.

Parameters

b
event
any

Returns

error
core/events/bus.go:333-341
func EmitAny(ctx context.Context, b *Bus, event any) error

{
	if b == nil {
		b = defaultBus
	}
	if event == nil {
		return nil
	}
	return emitByType(ctx, b, reflect.TypeOf(event), event)
}
F
function

emitByType

Parameters

b
event
any

Returns

error
core/events/bus.go:343-385
func emitByType(ctx context.Context, b *Bus, key reflect.Type, event any) error

{
	subs, ok := b.subscribers.Get(key)

	b.mu.RLock()
	mws := b.middlewares
	b.mu.RUnlock()

	emit := func(ctx context.Context, evt any) error {
		if !ok {
			return nil
		}
		var errs []error
		for _, sub := range subs {
			fn := reflect.ValueOf(sub.handler)
			results := fn.Call([]reflect.Value{reflect.ValueOf(ctx), reflect.ValueOf(evt)})
			if len(results) == 1 && !results[0].IsNil() {
				err := results[0].Interface().(error)
				if b.strategy == StopOnFirstError {
					return err
				}
				errs = append(errs, err)
			}
		}
		if len(errs) > 0 {
			return errors.Join(errs...)
		}
		return nil
	}

	if len(mws) > 0 {
		chain := applyMiddleware(emit, mws)
		if err := chain(ctx, event); err != nil {
			return err
		}
		return emitWildcards(ctx, b, event)
	}

	if err := emit(ctx, event); err != nil {
		return err
	}

	return emitWildcards(ctx, b, event)
}
F
function

emitWildcards

Parameters

b
event
any

Returns

error
core/events/bus.go:387-400
func emitWildcards(ctx context.Context, b *Bus, event any) error

{
	b.mu.RLock()
	wildcards := b.wildcard
	b.mu.RUnlock()
	for _, w := range wildcards {
		if fn, ok := w.handler.(func(ctx context.Context, event any) error); ok {
			if err := fn(ctx, event); err != nil {
				return err
			}
		}
	}

	return nil
}
F
function

EmitAsync

EmitAsync dispatches an event asynchronously to the bus channel.

Parameters

b
event
T

Returns

error
core/events/bus.go:403-416
func EmitAsync[T any](ctx context.Context, b *Bus, event T) error

{
	if b == nil {
		b = defaultBus
	}
	evt := asyncEvent{
		event: event,
		ctx:   asyncContext(ctx),
		emit: func(ctx context.Context, evt any) error {
			return Emit(ctx, b, evt.(T))
		},
	}

	return emitAsyncEvent(b, evt)
}
F
function

EmitAnyAsync

EmitAnyAsync dispatches an event asynchronously using its runtime type.

Parameters

b
event
any

Returns

error
core/events/bus.go:419-427
func EmitAnyAsync(ctx context.Context, b *Bus, event any) error

{
	if b == nil {
		b = defaultBus
	}
	if event == nil {
		return nil
	}
	return emitAsyncByType(ctx, b, reflect.TypeOf(event), event)
}
F
function

emitAsyncByType

Parameters

b
event
any

Returns

error
core/events/bus.go:429-439
func emitAsyncByType(ctx context.Context, b *Bus, key reflect.Type, event any) error

{
	evt := asyncEvent{
		event: event,
		ctx:   asyncContext(ctx),
		emit: func(ctx context.Context, evt any) error {
			return emitByType(ctx, b, key, evt)
		},
	}

	return emitAsyncEvent(b, evt)
}
F
function

asyncContext

Parameters

Returns

core/events/bus.go:441-446
func asyncContext(ctx context.Context) context.Context

{
	if ctx == nil {
		return context.Background()
	}
	return ctx
}
F
function

emitAsyncEvent

Parameters

b

Returns

error
core/events/bus.go:448-484
func emitAsyncEvent(b *Bus, evt asyncEvent) error

{
	b.asyncMu.Lock()
	if b.asyncClosed {
		b.asyncMu.Unlock()
		return ErrBusClosed
	}
	b.asyncWG.Add(1)
	b.asyncMu.Unlock()
	defer b.asyncWG.Done()

	switch b.overflowStrat {
	case OverflowDropOldest:
		select {
		case b.asyncCh <- evt:
			return nil
		default:
			select {
			case <-b.asyncCh:
			default:
			}
			select {
			case b.asyncCh <- evt:
				return nil
			default:
				return ErrAsyncQueueFull
			}
		}
	case OverflowFail:
		select {
		case b.asyncCh <- evt:
			return nil
		default:
			return ErrAsyncQueueFull
		}
	}
	return ErrAsyncQueueFull
}
F
function

applyMiddleware

Parameters

handler
func(ctx context.Context, evt any) error
middlewares
core/events/bus.go:486-495
func applyMiddleware(handler func(ctx context.Context, evt any) error, middlewares []Middleware) func(ctx context.Context, evt any) error

{
	for i := len(middlewares) - 1; i >= 0; i-- {
		mw := middlewares[i]
		next := handler
		handler = func(ctx context.Context, evt any) error {
			return mw(ctx, evt, next)
		}
	}
	return handler
}
F
function

TestBus_SubscribeAndEmit

Parameters

core/events/bus_test.go:13-29
func TestBus_SubscribeAndEmit(t *testing.T)

{
	b := New()

	var received MyEvent
	Subscribe(b, func(ctx context.Context, e MyEvent) error {
		received = e
		return nil
	})

	err := Emit(context.Background(), b, MyEvent{ID: 1})
	if err != nil {
		t.Fatalf("Emit: %v", err)
	}
	if received.ID != 1 {
		t.Errorf("received %d, want 1", received.ID)
	}
}
S
struct

MyEvent

core/events/bus_test.go:31-31
type MyEvent struct

Fields

Name Type Description
ID int
F
function

TestBus_DefaultBus

Parameters

core/events/bus_test.go:33-46
func TestBus_DefaultBus(t *testing.T)

{
	defaultBus = New()

	var received MyEvent
	Subscribe(Default(), func(ctx context.Context, e MyEvent) error {
		received = e
		return nil
	})

	Emit(context.Background(), Default(), MyEvent{ID: 2})
	if received.ID != 2 {
		t.Errorf("received %d, want 2", received.ID)
	}
}
F
function

TestBus_NilBusDefaults

Parameters

core/events/bus_test.go:48-62
func TestBus_NilBusDefaults(t *testing.T)

{
	b := New()
	defer b.Close()

	var received MyEvent
	Subscribe(nil, func(ctx context.Context, e MyEvent) error {
		received = e
		return nil
	})

	Emit(context.Background(), nil, MyEvent{ID: 8})
	if received.ID != 8 {
		t.Errorf("received %d, want 8", received.ID)
	}
}
F
function

TestBus_EmitAsync

Parameters

core/events/bus_test.go:64-97
func TestBus_EmitAsync(t *testing.T)

{
	b := New()
	defer b.Close()

	var mu sync.Mutex
	var received []MyEvent
	done := make(chan struct{})

	Subscribe(b, func(ctx context.Context, e MyEvent) error {
		mu.Lock()
		defer mu.Unlock()
		received = append(received, e)
		if len(received) == 5 {
			close(done)
		}
		return nil
	})

	for i := 0; i < 5; i++ {
		EmitAsync(context.Background(), b, MyEvent{ID: i})
	}

	select {
	case <-done:
	case <-b.asyncClose:
		t.Fatal("bus closed before events processed")
	}

	mu.Lock()
	defer mu.Unlock()
	if len(received) != 5 {
		t.Fatalf("got %d events, want 5", len(received))
	}
}
F
function

TestBus_EmitAny

Parameters

core/events/bus_test.go:99-116
func TestBus_EmitAny(t *testing.T)

{
	b := New()
	defer b.Close()

	var received *MyEvent
	Subscribe(b, func(ctx context.Context, e *MyEvent) error {
		received = e
		return nil
	})

	err := EmitAny(context.Background(), b, &MyEvent{ID: 9})
	if err != nil {
		t.Fatalf("EmitAny: %v", err)
	}
	if received == nil || received.ID != 9 {
		t.Fatalf("received = %#v", received)
	}
}
F
function

TestBus_EmitAnyWithMiddlewareAndWildcard

Parameters

core/events/bus_test.go:118-139
func TestBus_EmitAnyWithMiddlewareAndWildcard(t *testing.T)

{
	b := New()
	defer b.Close()

	b.Use(func(ctx context.Context, event any, next func(context.Context, any) error) error {
		return next(ctx, event)
	})

	called := false
	SubscribeWildcard(b, func(ctx context.Context, event any) error {
		called = true
		return nil
	})

	err := EmitAny(context.Background(), b, &MyEvent{ID: 10})
	if err != nil {
		t.Fatalf("EmitAny: %v", err)
	}
	if !called {
		t.Fatal("wildcard handler was not called")
	}
}
F
function

TestBus_Middleware

Parameters

core/events/bus_test.go:141-163
func TestBus_Middleware(t *testing.T)

{
	b := New()
	defer b.Close()

	var chain []string
	b.Use(func(ctx context.Context, event any, next func(ctx context.Context, event any) error) error {
		chain = append(chain, "before")
		err := next(ctx, event)
		chain = append(chain, "after")
		return err
	})

	Subscribe(b, func(ctx context.Context, e MyEvent) error {
		chain = append(chain, "handler")
		return nil
	})

	Emit(context.Background(), b, MyEvent{ID: 3})

	if len(chain) != 3 || chain[0] != "before" || chain[1] != "handler" || chain[2] != "after" {
		t.Errorf("middleware chain wrong: %v", chain)
	}
}
F
function

TestBus_StrategyStopOnFirstError

Parameters

core/events/bus_test.go:165-185
func TestBus_StrategyStopOnFirstError(t *testing.T)

{
	b := New()
	defer b.Close()

	Subscribe(b, func(ctx context.Context, e MyEvent) error {
		return context.Canceled
	})
	called := false
	Subscribe(b, func(ctx context.Context, e MyEvent) error {
		called = true
		return nil
	})

	err := Emit(context.Background(), b, MyEvent{ID: 5})
	if err != context.Canceled {
		t.Errorf("expected context.Canceled, got %v", err)
	}
	if called {
		t.Error("second handler should not be called under StopOnFirstError")
	}
}
F
function

TestBus_OnAsyncError

Parameters

core/events/bus_test.go:187-203
func TestBus_OnAsyncError(t *testing.T)

{
	b := New(
		WithStrategy(StopOnFirstError),
		WithOnAsyncError(func(err error) {
			if err == nil {
				t.Error("expected error in async handler")
			}
		}),
	)
	defer b.Close()

	Subscribe(b, func(ctx context.Context, e MyEvent) error {
		return context.Canceled
	})

	EmitAsync(context.Background(), b, MyEvent{ID: 5})
}
F
function

TestBus_EmitAsyncUsesMiddlewareAndWildcard

Parameters

core/events/bus_test.go:205-235
func TestBus_EmitAsyncUsesMiddlewareAndWildcard(t *testing.T)

{
	b := New()
	defer b.Close()

	middlewareCalled := make(chan struct{}, 1)
	wildcardCalled := make(chan struct{}, 1)
	b.Use(func(ctx context.Context, event any, next func(context.Context, any) error) error {
		middlewareCalled <- struct{}{}
		return next(ctx, event)
	})
	SubscribeWildcard(b, func(context.Context, any) error {
		wildcardCalled <- struct{}{}
		return nil
	})
	Subscribe(b, func(context.Context, MyEvent) error {
		return nil
	})

	EmitAsync(context.Background(), b, MyEvent{ID: 6})

	for name, called := range map[string]<-chan struct{}{
		"middleware": middlewareCalled,
		"wildcard":   wildcardCalled,
	} {
		select {
		case <-called:
		case <-time.After(time.Second):
			t.Fatalf("timed out waiting for async %s", name)
		}
	}
}
F
function

TestBus_EmitAnyAsyncReportsBestEffortErrors

Parameters

core/events/bus_test.go:237-266
func TestBus_EmitAnyAsyncReportsBestEffortErrors(t *testing.T)

{
	firstErr := errors.New("first")
	secondErr := errors.New("second")
	asyncErr := make(chan error, 1)
	b := New(
		WithStrategy(BestEffort),
		WithOnAsyncError(func(err error) {
			asyncErr <- err
		}),
	)
	defer b.Close()

	Subscribe(b, func(context.Context, MyEvent) error {
		return firstErr
	})
	Subscribe(b, func(context.Context, MyEvent) error {
		return secondErr
	})

	EmitAnyAsync(context.Background(), b, MyEvent{ID: 7})

	select {
	case err := <-asyncErr:
		if !errors.Is(err, firstErr) || !errors.Is(err, secondErr) {
			t.Fatalf("async error = %v, want both handler errors", err)
		}
	case <-time.After(time.Second):
		t.Fatal("timed out waiting for async error")
	}
}
F
function

TestBus_Close

Parameters

core/events/bus_test.go:268-281
func TestBus_Close(t *testing.T)

{
	b := New()
	b.Close()
	b.Close()
	if err := b.Shutdown(context.Background()); err != nil {
		t.Fatal(err)
	}
	done := make(chan struct{})
	go func() {
		EmitAsync(context.Background(), b, MyEvent{ID: 1})
		close(done)
	}()
	<-done
}
F
function

TestBus_CloseDrainsQueuedEvents

Parameters

core/events/bus_test.go:283-306
func TestBus_CloseDrainsQueuedEvents(t *testing.T)

{
	b := New(WithBufferSize(16))
	var handled int
	var mu sync.Mutex
	Subscribe(b, func(context.Context, MyEvent) error {
		mu.Lock()
		handled++
		mu.Unlock()
		return nil
	})
	for index := 0; index < 10; index++ {
		EmitAsync(context.Background(), b, MyEvent{ID: index})
	}

	if err := b.Shutdown(context.Background()); err != nil {
		t.Fatal(err)
	}

	mu.Lock()
	defer mu.Unlock()
	if handled != 10 {
		t.Fatalf("Shutdown() drained %d events, want 10", handled)
	}
}
F
function

TestBus_AsyncHandlerCanCloseBus

Parameters

core/events/bus_test.go:308-328
func TestBus_AsyncHandlerCanCloseBus(t *testing.T)

{
	b := New(WithBufferSize(1))
	handlerDone := make(chan struct{})
	Subscribe(b, func(context.Context, MyEvent) error {
		b.Close()
		close(handlerDone)
		return nil
	})
	EmitAsync(context.Background(), b, MyEvent{})

	select {
	case <-handlerDone:
	case <-time.After(time.Second):
		t.Fatal("async handler deadlocked in Close()")
	}
	select {
	case <-b.shutdownDone:
	case <-time.After(time.Second):
		t.Fatal("bus did not finish shutdown after reentrant Close()")
	}
}
F
function

TestBus_AsyncHandlerCannotWaitForOwnShutdown

Parameters

core/events/bus_test.go:330-352
func TestBus_AsyncHandlerCannotWaitForOwnShutdown(t *testing.T)

{
	b := New()
	handlerDone := make(chan error, 1)
	Subscribe(b, func(ctx context.Context, _ MyEvent) error {
		handlerDone <- b.Shutdown(ctx)
		return nil
	})
	EmitAsync(context.Background(), b, MyEvent{})

	select {
	case err := <-handlerDone:
		if !errors.Is(err, ErrReentrantShutdown) {
			t.Fatalf("Shutdown() error = %v, want ErrReentrantShutdown", err)
		}
	case <-time.After(time.Second):
		t.Fatal("async handler deadlocked in Shutdown()")
	}
	select {
	case <-b.shutdownDone:
	case <-time.After(time.Second):
		t.Fatal("bus did not finish after reentrant Shutdown()")
	}
}
F
function

TestBus_AsyncHandlerCanEmitAsync

Parameters

core/events/bus_test.go:354-376
func TestBus_AsyncHandlerCanEmitAsync(t *testing.T)

{
	emitErr := make(chan error, 1)
	b := New(WithBufferSize(1))
	Subscribe(b, func(ctx context.Context, _ MyEvent) error {
		emitErr <- EmitAsync(ctx, b, 1)
		return nil
	})
	if err := EmitAsync(context.Background(), b, MyEvent{}); err != nil {
		t.Fatalf("EmitAsync: %v", err)
	}

	select {
	case err := <-emitErr:
		if err != nil {
			t.Fatalf("reentrant emission error = %v", err)
		}
	case <-time.After(time.Second):
		t.Fatal("async handler deadlocked while emitting")
	}
	if err := b.Shutdown(context.Background()); err != nil {
		t.Fatal(err)
	}
}
F
function

TestBus_DropOldestDoesNotBlockReentrantEmission

Parameters

core/events/bus_test.go:378-407
func TestBus_DropOldestDoesNotBlockReentrantEmission(t *testing.T)

{
	emitErr := make(chan error, 1)
	b := New(
		WithBufferSize(1),
		WithOverflowStrategy(OverflowDropOldest),
	)
	Subscribe(b, func(ctx context.Context, _ MyEvent) error {
		if err := EmitAsync(ctx, b, 1); err != nil {
			emitErr <- err
			return nil
		}
		emitErr <- EmitAsync(ctx, b, 2)
		return nil
	})
	if err := EmitAsync(context.Background(), b, MyEvent{}); err != nil {
		t.Fatalf("EmitAsync: %v", err)
	}

	select {
	case err := <-emitErr:
		if err != nil {
			t.Fatalf("reentrant emission error = %v", err)
		}
	case <-time.After(time.Second):
		t.Fatal("DropOldest deadlocked on a reentrant unbuffered emission")
	}
	if err := b.Shutdown(context.Background()); err != nil {
		t.Fatal(err)
	}
}
F
function

TestBus_AsyncErrorCallbackCanEmitToSameBus

Parameters

core/events/bus_test.go:409-440
func TestBus_AsyncErrorCallbackCanEmitToSameBus(t *testing.T)

{
	emitErr := make(chan error, 1)
	var b *Bus
	b = New(
		WithBufferSize(1),
		WithOnAsyncError(func(error) {
			if err := EmitAsync(context.Background(), b, 1); err != nil {
				emitErr <- err
				return
			}
			emitErr <- EmitAsync(context.Background(), b, 2)
		}),
	)
	Subscribe(b, func(context.Context, MyEvent) error {
		return errors.New("handler failed")
	})
	if err := EmitAsync(context.Background(), b, MyEvent{}); err != nil {
		t.Fatalf("EmitAsync: %v", err)
	}

	select {
	case err := <-emitErr:
		if !errors.Is(err, ErrAsyncQueueFull) {
			t.Fatalf("callback emission error = %v, want ErrAsyncQueueFull", err)
		}
	case <-time.After(time.Second):
		t.Fatal("async error callback deadlocked while emitting to the same bus")
	}
	if err := b.Shutdown(context.Background()); err != nil {
		t.Fatal(err)
	}
}
F
function

TestBus_AsyncPanicBecomesError

Parameters

core/events/bus_test.go:442-462
func TestBus_AsyncPanicBecomesError(t *testing.T)

{
	asyncErr := make(chan error, 1)
	b := New(WithOnAsyncError(func(err error) {
		asyncErr <- err
	}))
	defer b.Close()
	Subscribe(b, func(context.Context, MyEvent) error {
		panic("handler failed")
	})

	EmitAsync(context.Background(), b, MyEvent{})

	select {
	case err := <-asyncErr:
		if err == nil || !strings.Contains(err.Error(), "handler failed") {
			t.Fatalf("async panic error = %v", err)
		}
	case <-time.After(time.Second):
		t.Fatal("timed out waiting for recovered panic")
	}
}
F
function

TestBus_CustomBufferProcessesEvents

Parameters

core/events/bus_test.go:464-474
func TestBus_CustomBufferProcessesEvents(t *testing.T)

{
	b := New(WithBufferSize(1))
	defer b.Close()
	received := make(chan struct{}, 1)
	Subscribe(b, func(context.Context, MyEvent) error {
		received <- struct{}{}
		return nil
	})
	EmitAsync(context.Background(), b, MyEvent{ID: 1})
	<-received
}
F
function

TestBus_EmitAsyncPreservesContext

Parameters

core/events/bus_test.go:476-497
func TestBus_EmitAsyncPreservesContext(t *testing.T)

{
	type contextKey struct{}
	bus := New()
	defer bus.Close()
	received := make(chan string, 1)
	Subscribe(bus, func(ctx context.Context, event int) error {
		value, _ := ctx.Value(contextKey{}).(string)
		received <- value
		return nil
	})
	ctx := context.WithValue(context.Background(), contextKey{}, "tenant")
	EmitAsync(ctx, bus, 1)

	select {
	case value := <-received:
		if value != "tenant" {
			t.Fatalf("context value = %q, want tenant", value)
		}
	case <-time.After(time.Second):
		t.Fatal("timed out waiting for async event")
	}
}
F
function

TestBus_RejectsInvalidBufferSize

Parameters

core/events/bus_test.go:499-510
func TestBus_RejectsInvalidBufferSize(t *testing.T)

{
	for _, size := range []int{-1, 0} {
		t.Run(fmt.Sprintf("size_%d", size), func(t *testing.T) {
			defer func() {
				if recover() == nil {
					t.Fatal("New() accepted an invalid async buffer")
				}
			}()
			New(WithBufferSize(size))
		})
	}
}
F
function

TestBus_EmitNoSubscribers

Parameters

core/events/bus_test.go:512-519
func TestBus_EmitNoSubscribers(t *testing.T)

{
	b := New()
	defer b.Close()
	err := Emit(context.Background(), b, MyEvent{ID: 7})
	if err != nil {
		t.Errorf("Emit without subscribers: %v", err)
	}
}
F
function

TestBus_NilBusDefault

Parameters

core/events/bus_test.go:521-530
func TestBus_NilBusDefault(t *testing.T)

{
	var b *Bus = nil
	Subscribe(b, func(ctx context.Context, e MyEvent) error {
		return nil
	})
	err := Emit(context.Background(), b, MyEvent{ID: 8})
	if err != nil {
		t.Errorf("Emit with nil bus: %v", err)
	}
}