Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 5 additions & 2 deletions aggregatestore/events/aggregatebase.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,9 +93,12 @@ func (a *AggregateBase) IncrementVersion() {

// Events implements the Events method of the eh.EventSource interface.
func (a *AggregateBase) Events() []eh.Event {
events := a.events
return a.events
}

// ClearEvents implements the ClearEvents method of the eh.EventSource interface.
func (a *AggregateBase) ClearEvents() {
a.events = nil
return events
}

// AppendEvent appends an event for later retrieval by Events().
Expand Down
6 changes: 4 additions & 2 deletions aggregatestore/events/aggregatebase_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ func TestAggregateEvents(t *testing.T) {
if event1.String() != "TestAggregateEvent@1" {
t.Error("the string representation should be correct:", event1.String())
}
events := agg.events
events := agg.Events()
if len(events) != 1 {
t.Fatal("there should be one event provided:", len(events))
}
Expand All @@ -95,6 +95,7 @@ func TestAggregateEvents(t *testing.T) {
if len(events) != 2 {
t.Error("there should be two events provided:", len(events))
}
agg.ClearEvents()

event3 := agg.AppendEvent(TestAggregateEventType, &TestEventData{"event1"}, timestamp)
if event3.Version() != 1 {
Expand All @@ -104,11 +105,12 @@ func TestAggregateEvents(t *testing.T) {
if len(events) != 1 {
t.Error("there should be one new event provided:", len(events))
}
agg.ClearEvents()

agg = NewTestAggregate(uuid.New())
event1 = agg.AppendEvent(TestAggregateEventType, &TestEventData{"event1"}, timestamp)
event2 = agg.AppendEvent(TestAggregateEventType, &TestEventData{"event2"}, timestamp)
events = agg.events
events = agg.Events()
if len(events) != 2 {
t.Fatal("there should be 2 events provided:", len(events))
}
Expand Down
1 change: 1 addition & 0 deletions aggregatestore/events/aggregatestore.go
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,7 @@ func (r *AggregateStore) Save(ctx context.Context, agg eh.Aggregate) error {
if err := r.store.Save(ctx, events, a.Version()); err != nil {
return err
}
a.ClearEvents()

// Apply the events in case the aggregate needs to be further used
// after this save. Currently it is not reused.
Expand Down
4 changes: 3 additions & 1 deletion aggregatestore/model/aggregatestore.go
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,9 @@ func (r *AggregateStore) Save(ctx context.Context, aggregate eh.Aggregate) error

// Handle any events optionally provided by the aggregate.
if a, ok := aggregate.(eh.EventSource); ok && r.eventHandler != nil {
for _, e := range a.Events() {
events := a.Events()
a.ClearEvents()
for _, e := range events {
if err := r.eventHandler.HandleEvent(ctx, e); err != nil {
return err
}
Expand Down
7 changes: 5 additions & 2 deletions aggregatestore/model/eventsource.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,10 @@ func (a *SliceEventSource) AppendEvent(e eh.Event) {

// Events implements the Events method of the eh.EventSource interface.
func (a *SliceEventSource) Events() []eh.Event {
events := *a
return *a
}

// ClearEvents implements the ClearEvents method of the eh.EventSource interface.
func (a *SliceEventSource) ClearEvents() {
*a = nil
return events
}
2 changes: 1 addition & 1 deletion docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ services:
- "27017:27017"

gpubsub:
image: gcr.io/google.com/cloudsdktool/cloud-sdk:326.0.0-emulators
image: gcr.io/google.com/cloudsdktool/cloud-sdk:329.0.0-emulators
ports:
- "8793:8793"
entrypoint:
Expand Down
28 changes: 27 additions & 1 deletion eventbus/acceptance_testing.go
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,7 @@ func AcceptanceTest(t *testing.T, bus1, bus2 eh.EventBus, timeout time.Duration)

// Test async errors from handlers.
errorHandler := mocks.NewEventHandler("error_handler")
errorHandler.Err = errors.New("handler error")
errorHandler.ErrOnce = errors.New("handler error")
if err := bus1.AddHandler(ctx, eh.MatchAll{}, errorHandler); err != nil {
t.Fatal("there should be no error:", err)
}
Expand All @@ -238,9 +238,35 @@ func AcceptanceTest(t *testing.T, bus1, bus2 eh.EventBus, timeout time.Duration)
// Good case.
if err.Error() != "could not handle event (error_handler): handler error: (Event@3)" {
t.Error("incorrect error sent on event bus:", err)
t.Logf("%#v", err.Event)
}
}

// Retryable events.
retryHandler := mocks.NewEventHandler("retry_handler")
retryHandler.ErrOnce = eh.RetryableEventError{Err: errors.New("retryable error")}
bus1.AddHandler(ctx, eh.MatchAll{}, retryHandler)

time.Sleep(timeout) // Need to wait here for handlers to be added.

event4 := eh.NewEvent(mocks.EventType, &mocks.EventData{Content: "event4"}, timestamp,
eh.ForAggregate(mocks.AggregateType, id, 4),
eh.WithMetadata(map[string]interface{}{"meta": "data", "num": int32(42)}),
)
if err := bus1.HandleEvent(ctx, event4); err != nil {
t.Error("there should be no error:", err)
}
select {
case <-time.After(timeout):
t.Error("there should be a retried event in time")
case <-retryHandler.Recv:
}
retryHandler.Lock()
if retryHandler.NumHandleEvent != 2 {
t.Error("the handler should have been called twice")
}
retryHandler.Unlock()

// Cancel all handlers and wait.
cancel()
bus1.Wait()
Expand Down
9 changes: 7 additions & 2 deletions eventbus/gcp/eventbus.go
Original file line number Diff line number Diff line change
Expand Up @@ -253,14 +253,19 @@ func (b *EventBus) handler(m eh.EventMatcher, h eh.EventHandler) func(ctx contex

// Handle the event if it did match.
if err := h.HandleEvent(ctx, event); err != nil {
// Retryable errors are not logged and will be retried.
if _, ok := err.(eh.RetryableEventError); ok {
msg.Nack()
return
}

// Log unhandled events, they will NOT be retried.
err = fmt.Errorf("could not handle event (%s): %w", h.HandlerType(), err)
select {
case b.errCh <- eh.EventBusError{Err: err, Ctx: ctx, Event: event}:
default:
log.Printf("eventhorizon: missed error in GCP event bus: %s", err)
}
msg.Nack()
return
}

msg.Ack()
Expand Down
9 changes: 7 additions & 2 deletions eventbus/jetstream/eventbus.go
Original file line number Diff line number Diff line change
Expand Up @@ -216,14 +216,19 @@ func (b *EventBus) handler(ctx context.Context, m eh.EventMatcher, h eh.EventHan

// Handle the event if it did match.
if err := h.HandleEvent(ctx, event); err != nil {
// Retryable errors are not logged and will be retried.
if _, ok := err.(eh.RetryableEventError); ok {
msg.Nak()
return
}

// Log unhandled events, they will NOT be retried.
err = fmt.Errorf("could not handle event (%s): %w", h.HandlerType(), err)
select {
case b.errCh <- eh.EventBusError{Err: err, Ctx: ctx, Event: event}:
default:
log.Printf("eventhorizon: missed error in Jetstream event bus: %s", err)
}
msg.Nak()
return
}

msg.AckSync()
Expand Down
7 changes: 6 additions & 1 deletion eventbus/kafka/eventbus.go
Original file line number Diff line number Diff line change
Expand Up @@ -267,13 +267,18 @@ func (b *EventBus) handler(m eh.EventMatcher, h eh.EventHandler, r *kafka.Reader

// Handle the event if it did match.
if err := h.HandleEvent(ctx, event); err != nil {
// Retryable errors are not logged and will be retried.
if _, ok := err.(eh.RetryableEventError); ok {
return
}

// Log unhandled events, they will NOT be retried.
err = fmt.Errorf("could not handle event (%s): %w", h.HandlerType(), err)
select {
case b.errCh <- eh.EventBusError{Err: err, Ctx: ctx, Event: event}:
default:
log.Printf("eventhorizon: missed error in Kafka event bus: %s", err)
}
return
}

r.CommitMessages(ctx, msg)
Expand Down
18 changes: 15 additions & 3 deletions eventbus/local/eventbus.go
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,7 @@ type evt struct {
}

// Handles all events coming in on the channel.
func (b *EventBus) handle(ctx context.Context, m eh.EventMatcher, h eh.EventHandler, ch <-chan []byte) {
func (b *EventBus) handle(ctx context.Context, m eh.EventMatcher, h eh.EventHandler, ch chan []byte) {
defer b.wg.Done()

for {
Expand All @@ -161,7 +161,19 @@ func (b *EventBus) handle(ctx context.Context, m eh.EventMatcher, h eh.EventHand

// Handle the event if it did match.
if err := h.HandleEvent(ctx, event); err != nil {
err = fmt.Errorf("could not handle event (%s): %s", h.HandlerType(), err.Error())
// Retryable errors are not logged and will be retried.
if _, ok := err.(eh.RetryableEventError); ok {
select {
case ch <- data:
// Retry event by putting it back on the bus.
default:
log.Printf("eventhorizon: publish queue full for retry in local event bus")
}
continue
}

// Log unhandled events, they will NOT be retried.
err = fmt.Errorf("could not handle event (%s): %w", h.HandlerType(), err)
select {
case b.errCh <- eh.EventBusError{Err: err, Ctx: ctx, Event: event}:
default:
Expand All @@ -187,7 +199,7 @@ func NewGroup() *Group {
}
}

func (g *Group) channel(id string) <-chan []byte {
func (g *Group) channel(id string) chan []byte {
g.busMu.Lock()
defer g.busMu.Unlock()

Expand Down
19 changes: 16 additions & 3 deletions eventbus/redis/eventbus.go
Original file line number Diff line number Diff line change
Expand Up @@ -190,11 +190,12 @@ func (b *EventBus) handle(ctx context.Context, m eh.EventMatcher, h eh.EventHand
defer b.wg.Done()

msgHandler := b.handler(m, h, groupName)
readOpt := ">"
for {
streams, err := b.client.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: groupName,
Consumer: groupName + "_" + b.clientID,
Streams: []string{b.streamName, ">"},
Streams: []string{b.streamName, readOpt},
}).Result()
if err != nil && err != context.Canceled {
err = fmt.Errorf("could not receive: %w", err)
Expand All @@ -219,6 +220,13 @@ func (b *EventBus) handle(ctx context.Context, m eh.EventMatcher, h eh.EventHand
msgHandler(ctx, &msg)
}
}

// Flip flop the read option to read new and non-acked messages every other time.
if readOpt == ">" {
readOpt = "0"
} else {
readOpt = ">"
}
}
}

Expand Down Expand Up @@ -253,14 +261,19 @@ func (b *EventBus) handler(m eh.EventMatcher, h eh.EventHandler, groupName strin

// Handle the event if it did match.
if err := h.HandleEvent(ctx, event); err != nil {
// Retryable errors are not logged and will be retried.
if _, ok := err.(eh.RetryableEventError); ok {
// TODO: Nack if possible.
return
}

// Log unhandled events, they will NOT be retried.
err = fmt.Errorf("could not handle event (%s): %w", h.HandlerType(), err)
select {
case b.errCh <- eh.EventBusError{Err: err, Ctx: ctx, Event: event}:
default:
log.Printf("eventhorizon: missed error in Redis event bus: %s", err)
}
// TODO: Nack if possible.
return
}

_, err = b.client.XAck(ctx, b.streamName, groupName, msg.ID).Result()
Expand Down
19 changes: 19 additions & 0 deletions eventhandler.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ package eventhorizon

import (
"context"
"fmt"
"reflect"
"runtime"
"strings"
Expand All @@ -40,6 +41,24 @@ type EventHandler interface {
HandleEvent(context.Context, Event) error
}

// RetryableEventError is a "soft" error that handlers should return if they want the
// handler to be retried. This will often be the case when handling events (for
// example in a saga) where related read models have not yet been projected.
// NOTE: The retry behavior is dependent on the eventbus implementation used.
type RetryableEventError struct {
Err error
}

// Error implements the Error method of the error interface.
func (e RetryableEventError) Error() string {
return fmt.Sprintf("retryable: %s", e.Err)
}

// Cause returns the cause of this error.
func (e RetryableEventError) Cause() error {
return e.Err
}

// EventHandlerFunc is a function that can be used as a event handler.
type EventHandlerFunc func(context.Context, Event) error

Expand Down
Loading