mock_event_consumer.go raw

   1  package tests
   2  
   3  import (
   4  	"context"
   5  	"sync"
   6  	"time"
   7  
   8  	"github.com/getAlby/hub/events"
   9  )
  10  
  11  type mockEventConsumer struct {
  12  	mtx            sync.Mutex
  13  	consumedEvents []*events.Event
  14  }
  15  
  16  func NewMockEventConsumer() *mockEventConsumer {
  17  	return &mockEventConsumer{
  18  		consumedEvents: []*events.Event{},
  19  	}
  20  }
  21  
  22  func (e *mockEventConsumer) ConsumeEvent(ctx context.Context, event *events.Event, globalProperties map[string]interface{}) {
  23  	e.mtx.Lock()
  24  	defer e.mtx.Unlock()
  25  	e.consumedEvents = append(e.consumedEvents, event)
  26  }
  27  
  28  func (e *mockEventConsumer) GetConsumedEvents() []*events.Event {
  29  	// events are consumed async - give it a bit of time for tests
  30  	time.Sleep(10 * time.Millisecond)
  31  	return e.snapshotConsumedEvents()
  32  }
  33  
  34  // WaitForConsumedEvents waits until at least count events have been consumed
  35  // (events are consumed async) and returns them. On timeout it returns the
  36  // events consumed so far, so the caller's assertions fail with a useful message.
  37  func (e *mockEventConsumer) WaitForConsumedEvents(count int) []*events.Event {
  38  	deadline := time.Now().Add(5 * time.Second)
  39  	for {
  40  		consumedEvents := e.snapshotConsumedEvents()
  41  		if len(consumedEvents) >= count || time.Now().After(deadline) {
  42  			return consumedEvents
  43  		}
  44  		time.Sleep(10 * time.Millisecond)
  45  	}
  46  }
  47  
  48  func (e *mockEventConsumer) snapshotConsumedEvents() []*events.Event {
  49  	e.mtx.Lock()
  50  	defer e.mtx.Unlock()
  51  	return append([]*events.Event{}, e.consumedEvents...)
  52  }
  53