Skip to content

Commit 06365d9

Browse files
committed
Handle events asynchronously
That should help to keep up with the stream of messages and prevent services from being flagged as slow consumers.
1 parent dbc8a21 commit 06365d9

3 files changed

Lines changed: 15 additions & 12 deletions

File tree

services/clientlog/pkg/service/service.go

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -83,7 +83,8 @@ EventLoop:
8383
if !ok {
8484
break EventLoop
8585
}
86-
cl.processEvent(event)
86+
87+
go cl.processEvent(event)
8788

8889
if cl.stopped.Load() {
8990
break EventLoop

services/sse/pkg/service/service.go

Lines changed: 12 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -49,17 +49,19 @@ func (s SSE) ServeHTTP(w http.ResponseWriter, r *http.Request) {
4949
// ListenForEvents listens for events
5050
func (s SSE) ListenForEvents() {
5151
for e := range s.evChannel {
52-
switch ev := e.Event.(type) {
53-
default:
54-
s.l.Error().Interface("event", ev).Msg("unhandled event")
55-
case events.SendSSE:
56-
for _, uid := range ev.UserIDs {
57-
s.sse.Publish(uid, &sse.Event{
58-
Event: []byte(ev.Type),
59-
Data: ev.Message,
60-
})
52+
go func() {
53+
switch ev := e.Event.(type) {
54+
default:
55+
s.l.Error().Interface("event", ev).Msg("unhandled event")
56+
case events.SendSSE:
57+
for _, uid := range ev.UserIDs {
58+
s.sse.Publish(uid, &sse.Event{
59+
Event: []byte(ev.Type),
60+
Data: ev.Message,
61+
})
62+
}
6163
}
62-
}
64+
}()
6365
}
6466
}
6567

services/userlog/pkg/service/service.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -102,7 +102,7 @@ func (ul *UserlogService) MemorizeEvents(ch <-chan events.Event) {
102102
for i := 0; i < ul.cfg.MaxConcurrency; i++ {
103103
go func(ch <-chan events.Event) {
104104
for event := range ch {
105-
ul.processEvent(event)
105+
go ul.processEvent(event)
106106
}
107107
}(ch)
108108
}

0 commit comments

Comments
 (0)