-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathinternal_message_handler.go
More file actions
58 lines (52 loc) · 1.74 KB
/
Copy pathinternal_message_handler.go
File metadata and controls
58 lines (52 loc) · 1.74 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
package gluon
import (
"context"
"reflect"
)
// InternalMessageHandler is the message handler used by concrete drivers.
type InternalMessageHandler func(ctx context.Context, subscriber *Subscriber, message *TransportMessage) error
func getInternalHandler(b *Bus) InternalMessageHandler {
return func(ctx context.Context, sub *Subscriber, msg *TransportMessage) error {
msgMeta := b.internalSchemaRegistry.getByTopic(sub.key)
data := reflect.New(msgMeta.SchemaInternalType)
var schemaDef string
var err error
if b.SchemaRegistry != nil {
schemaDef, err = b.SchemaRegistry.GetSchemaDefinition(msgMeta.SchemaName, msgMeta.SchemaVersion)
if err != nil {
return err
}
}
err = b.Marshaler.Unmarshal(schemaDef, msg.Data, data.Interface())
logInternalConsumerError(b, err)
if err != nil && b.isLoggerEnabled() {
return err
}
return execConsumer(ctx, b, sub, msg, data)
}
}
func logInternalConsumerError(b *Bus, err error) {
if err != nil && b.isLoggerEnabled() {
b.Logger.Print("gluon: " + err.Error())
}
}
func execConsumer(ctx context.Context, b *Bus, sub *Subscriber, msg *TransportMessage, data reflect.Value) error {
scopedCtx := injectCorrelationContext(ctx, msg)
handlerFunc := sub.GetDefaultHandler()
for _, mw := range b.consumerMiddleware {
if mw != nil {
handlerFunc = mw(handlerFunc)
}
}
return handlerFunc(scopedCtx, &Message{
Headers: generateHeaders(msg, sub),
Data: data.Elem().Interface(),
})
}
func injectCorrelationContext(ctx context.Context, msg *TransportMessage) context.Context {
if msg.CorrelationID != "" {
ctx = context.WithValue(ctx, contextCorrelationID, gluonContextKey(msg.CorrelationID))
}
ctx = context.WithValue(ctx, contextMessageID, gluonContextKey(msg.ID))
return ctx
}