events
packageAPI reference for the events
package.
Imports
(11)Handler
Handler processes an event of type T.
type Handler func(ctx context.Context, event T) error
Priority
Priority defines the ordering of event handlers.
type Priority int
DispatchStrategy
DispatchStrategy controls error handling during event dispatch.
type DispatchStrategy int
Middleware
Middleware wraps event dispatch with cross-cutting behavior.
type Middleware func(ctx context.Context, event any, next func(ctx context.Context, event any) error) error
asyncEvent
type asyncEvent struct
Fields
| Name | Type | Description |
|---|---|---|
| event | any | |
| ctx | context.Context | |
| emit | func(ctx context.Context, event any) error |
OverflowStrategy
OverflowStrategy controls behavior when the async channel is full.
type OverflowStrategy int
Bus
Bus is the event bus that dispatches events to registered handlers.
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)
}
}
}
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
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 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 drains queued events and waits for active handlers.
Parameters
Returns
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 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{} |
subscriber
type subscriber struct
Fields
| Name | Type | Description |
|---|---|---|
| handler | any | |
| priority | Priority |
Uses
asyncBusContextKey
type asyncBusContextKey struct
Default
Default returns the package-level default Bus.
Returns
func Default() *Bus
{
return defaultBus
}
Option
Option configures a Bus.
type Option options.Option[Bus]
New
New creates a new Bus with the given options.
Parameters
Returns
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
}
WithBufferSize
WithBufferSize sets the async channel buffer size.
Parameters
Returns
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)
}
}
Uses
WithOverflowStrategy
WithOverflowStrategy sets the overflow behavior for async events.
Parameters
Returns
func WithOverflowStrategy(s OverflowStrategy) Option
{
return func(b *Bus) {
if s != OverflowFail && s != OverflowDropOldest {
panic("events: invalid overflow strategy")
}
b.overflowStrat = s
}
}
WithStrategy
WithStrategy sets the dispatch strategy for the Bus.
Parameters
Returns
func WithStrategy(s DispatchStrategy) Option
{
return func(b *Bus) { b.strategy = s }
}
WithOnAsyncError
WithOnAsyncError sets the error handler for async event emissions.
Parameters
Returns
func WithOnAsyncError(fn func(error)) Option
{
return func(b *Bus) { b.onAsyncError = fn }
}
Uses
Subscribe
Subscribe registers a typed handler for events of type T.
Parameters
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
})
}
SubscribeWildcard
SubscribeWildcard registers a handler for all event types.
Parameters
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})
}
Emit
Emit dispatches an event to all matching handlers synchronously.
Parameters
Returns
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
}
EmitAny
EmitAny dispatches an event using its runtime type.
Parameters
Returns
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)
}
emitByType
Parameters
Returns
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)
}
emitWildcards
Parameters
Returns
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
}
EmitAsync
EmitAsync dispatches an event asynchronously to the bus channel.
Parameters
Returns
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)
}
EmitAnyAsync
EmitAnyAsync dispatches an event asynchronously using its runtime type.
Parameters
Returns
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)
}
emitAsyncByType
Parameters
Returns
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)
}
asyncContext
Parameters
Returns
func asyncContext(ctx context.Context) context.Context
{
if ctx == nil {
return context.Background()
}
return ctx
}
emitAsyncEvent
Parameters
Returns
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
}
Uses
applyMiddleware
Parameters
Returns
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
}
TestBus_SubscribeAndEmit
Parameters
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)
}
}
MyEvent
type MyEvent struct
Fields
| Name | Type | Description |
|---|---|---|
| ID | int |
TestBus_DefaultBus
Parameters
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)
}
}
TestBus_NilBusDefaults
Parameters
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)
}
}
TestBus_EmitAsync
Parameters
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))
}
}
TestBus_EmitAny
Parameters
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)
}
}
TestBus_EmitAnyWithMiddlewareAndWildcard
Parameters
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")
}
}
TestBus_Middleware
Parameters
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)
}
}
TestBus_StrategyStopOnFirstError
Parameters
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")
}
}
TestBus_OnAsyncError
Parameters
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})
}
TestBus_EmitAsyncUsesMiddlewareAndWildcard
Parameters
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)
}
}
}
TestBus_EmitAnyAsyncReportsBestEffortErrors
Parameters
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")
}
}
TestBus_Close
Parameters
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
}
TestBus_CloseDrainsQueuedEvents
Parameters
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)
}
}
TestBus_AsyncHandlerCanCloseBus
Parameters
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()")
}
}
TestBus_AsyncHandlerCannotWaitForOwnShutdown
Parameters
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()")
}
}
TestBus_AsyncHandlerCanEmitAsync
Parameters
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)
}
}
TestBus_DropOldestDoesNotBlockReentrantEmission
Parameters
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)
}
}
TestBus_AsyncErrorCallbackCanEmitToSameBus
Parameters
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)
}
}
TestBus_AsyncPanicBecomesError
Parameters
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")
}
}
TestBus_CustomBufferProcessesEvents
Parameters
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
}
TestBus_EmitAsyncPreservesContext
Parameters
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")
}
}
TestBus_RejectsInvalidBufferSize
Parameters
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))
})
}
}
TestBus_EmitNoSubscribers
Parameters
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)
}
}
TestBus_NilBusDefault
Parameters
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)
}
}