This repository has been archived by the owner on Jul 28, 2020. It is now read-only.
forked from ThreeDotsLabs/watermill-kafka
-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathmarshaler_test.go
118 lines (89 loc) · 2.77 KB
/
marshaler_test.go
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
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
package kafka_test
import (
"testing"
"github.com/Shopify/sarama"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/ThreeDotsLabs/watermill"
"github.com/ThreeDotsLabs/watermill-kafka/v2/pkg/kafka"
"github.com/ThreeDotsLabs/watermill/message"
)
func TestDefaultMarshaler_MarshalUnmarshal(t *testing.T) {
m := kafka.DefaultMarshaler{}
msg := message.NewMessage(watermill.NewUUID(), []byte("payload"))
msg.Metadata.Set("foo", "bar")
marshaled, err := m.Marshal("topic", msg)
require.NoError(t, err)
unmarshaledMsg, err := m.Unmarshal(producerToConsumerMessage(marshaled))
require.NoError(t, err)
assert.True(t, msg.Equals(unmarshaledMsg))
}
func BenchmarkDefaultMarshaler_Marshal(b *testing.B) {
m := kafka.DefaultMarshaler{}
msg := message.NewMessage(watermill.NewUUID(), []byte("payload"))
msg.Metadata.Set("foo", "bar")
for i := 0; i < b.N; i++ {
m.Marshal("foo", msg)
}
}
func BenchmarkDefaultMarshaler_Unmarshal(b *testing.B) {
m := kafka.DefaultMarshaler{}
msg := message.NewMessage(watermill.NewUUID(), []byte("payload"))
msg.Metadata.Set("foo", "bar")
marshaled, err := m.Marshal("foo", msg)
if err != nil {
b.Fatal(err)
}
consumedMsg := producerToConsumerMessage(marshaled)
for i := 0; i < b.N; i++ {
m.Unmarshal(consumedMsg)
}
}
func TestWithPartitioningMarshaler_MarshalUnmarshal(t *testing.T) {
m := kafka.NewWithPartitioningMarshaler(func(topic string, msg *message.Message) (string, error) {
return msg.Metadata.Get("partition"), nil
})
partitionKey := "1"
msg := message.NewMessage(watermill.NewUUID(), []byte("payload"))
msg.Metadata.Set("partition", partitionKey)
producerMsg, err := m.Marshal("topic", msg)
require.NoError(t, err)
unmarshaledMsg, err := m.Unmarshal(producerToConsumerMessage(producerMsg))
require.NoError(t, err)
assert.True(t, msg.Equals(unmarshaledMsg))
assert.NoError(t, err)
producerKey, err := producerMsg.Key.Encode()
require.NoError(t, err)
assert.Equal(t, string(producerKey), partitionKey)
}
func producerToConsumerMessage(producerMessage *sarama.ProducerMessage) *sarama.ConsumerMessage {
var key []byte
if producerMessage.Key != nil {
var err error
key, err = producerMessage.Key.Encode()
if err != nil {
panic(err)
}
}
var value []byte
if producerMessage.Value != nil {
var err error
value, err = producerMessage.Value.Encode()
if err != nil {
panic(err)
}
}
var headers []*sarama.RecordHeader
for i := range producerMessage.Headers {
headers = append(headers, &producerMessage.Headers[i])
}
return &sarama.ConsumerMessage{
Key: key,
Value: value,
Topic: producerMessage.Topic,
Partition: producerMessage.Partition,
Offset: producerMessage.Offset,
Timestamp: producerMessage.Timestamp,
Headers: headers,
}
}