Concurrent processing of all async event handlers registered for the event type
(ctx context.Context, event Event)
| 102 | |
| 103 | // Concurrent processing of all async event handlers registered for the event type |
| 104 | func (em *EventManager) ProcessEvent(ctx context.Context, event Event) { |
| 105 | defer func() { <-em.activeCountChannel }() |
| 106 | // Send event to all registered handlers concurrently. WaitGroup blocks |
| 107 | // until all are finished |
| 108 | var wg sync.WaitGroup |
| 109 | for _, handler := range em.eventHandlers[event.EventType()] { |
| 110 | base.DebugfCtx(ctx, base.KeyEvents, "Event queue worker sending event %s to: %s", base.UD(event.String()), handler) |
| 111 | wg.Add(1) |
| 112 | go func(event Event, handler EventHandler) { |
| 113 | defer wg.Done() |
| 114 | //TODO: Currently we're not tracking success/fail from event handlers. When this |
| 115 | // is needed, could pass a channel to HandleEvent for tracking results |
| 116 | if handler.HandleEvent(ctx, event) { |
| 117 | em.IncrementEventsProcessedSuccess(1) |
| 118 | } else { |
| 119 | em.IncrementEventsProcessedFail(1) |
| 120 | } |
| 121 | base.TracefCtx(ctx, base.KeyAll, "Webhook event processed %s", event) |
| 122 | |
| 123 | }(event, handler) |
| 124 | } |
| 125 | wg.Wait() |
| 126 | } |
| 127 | |
| 128 | // Register a new event handler to the EventManager. The event manager will route events of |
| 129 | // type eventType to the handler. |
no test coverage detected