-
Notifications
You must be signed in to change notification settings - Fork 2
/
utils.go
85 lines (77 loc) · 2.8 KB
/
utils.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
package test
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/apache/pulsar-client-go/pulsar"
"github.com/kubescape/messaging/pulsar/connector"
)
func ProduceToTopic(suite *PulsarTestSuite, topic connector.TopicName, payloads [][]byte) {
producer, err := suite.Client.NewProducer(connector.WithProducerTopic(topic))
if err != nil {
suite.FailNow(err.Error(), "create producer")
}
defer producer.Close()
for _, payload := range payloads {
if _, err := producer.Send(context.Background(), &pulsar.ProducerMessage{Payload: payload}); err != nil {
suite.FailNow(err.Error(), "send payload")
}
}
}
func ProduceObjectsToTopic[P any](suite *PulsarTestSuite, topic connector.TopicName, payloads []P) {
producer, err := suite.Client.NewProducer(connector.WithProducerTopic(topic))
if err != nil {
suite.FailNow(err.Error(), "create producer")
}
defer producer.Close()
ProduceMessages[P](suite, context.Background(), producer, payloads)
}
func ProduceMessages[P any](suite *PulsarTestSuite, ctx context.Context, producer pulsar.Producer, payloads []P) {
for _, payload := range payloads {
payloadBytes, err := json.Marshal(payload)
if err != nil {
suite.FailNow(err.Error(), "marshal payload")
}
if _, err := producer.Send(ctx, &pulsar.ProducerMessage{Payload: payloadBytes}); err != nil {
suite.FailNow(err.Error(), "send payload")
}
}
}
// SubscribeToTopic - subscribe to a topic and returns a function that consumes messages from the topic
func SubscribeToTopic[P any](suite *PulsarTestSuite, topic connector.TopicName, payloads []P, subscription string) (consumeFunc func() []P) {
consumer, err := suite.Client.NewConsumer(connector.WithTopic(topic), connector.WithSubscriptionName(subscription))
if err != nil {
suite.FailNow(err.Error(), "subscribe")
}
return func() []P {
defer consumer.Close()
return ConsumeMessages[P](suite, context.Background(), consumer, subscription, 1)
}
}
func ConsumeMessages[P any](suite *PulsarTestSuite, ctx context.Context, consumer pulsar.Consumer, consumerId string, timeoutSeconds int) (actualPayloads []P) {
//consume payloads for X seconds
testConsumerCtx, consumerCancel := context.WithTimeout(ctx, time.Second*time.Duration(timeoutSeconds))
defer consumerCancel()
actualPayloads = []P{}
for {
msg, err := consumer.Receive(testConsumerCtx)
if err != nil {
if testConsumerCtx.Err() == nil {
fmt.Printf("%s: consumer error: %s", consumerId, err.Error())
suite.FailNow(err.Error(), "consumer error")
}
fmt.Printf("%s: breaking - %s", consumerId, err.Error())
break
}
var payload P
if err := json.Unmarshal(msg.Payload(), &payload); err != nil {
suite.NoError(err, "unmarshal failed")
consumer.Nack(msg)
continue
}
actualPayloads = append(actualPayloads, payload)
consumer.Ack(msg)
}
return actualPayloads
}