File size: 1,781 Bytes
d6f631f | 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 | package sink_test
import (
"testing"
"github.com/confluentinc/confluent-kafka-go/v2/kafka"
"github.com/stretchr/testify/assert"
"github.com/openmeterio/openmeter/openmeter/sink"
sinkmodels "github.com/openmeterio/openmeter/openmeter/sink/models"
)
func TestBuffer(t *testing.T) {
buffer := sink.NewSinkBuffer()
topic := "my-topic"
sinkMessage1 := sinkmodels.SinkMessage{
KafkaMessage: &kafka.Message{
TopicPartition: kafka.TopicPartition{
Topic: &topic,
Partition: 1,
Offset: 1,
},
},
}
sinkMessage2 := sinkmodels.SinkMessage{
KafkaMessage: &kafka.Message{
TopicPartition: kafka.TopicPartition{
Topic: &topic,
Partition: 1,
Offset: 2,
},
},
}
// We call add with the same message twice but as it has the
// same topic, partition and offset it should only be present in the buffer once.
buffer.Add(sinkMessage1)
buffer.Add(sinkMessage1)
buffer.Add(sinkMessage2)
assert.Equal(t, 2, buffer.Size())
assert.ElementsMatch(t, []sinkmodels.SinkMessage{sinkMessage1, sinkMessage2}, buffer.Dequeue())
}
func TestBufferRemoveByPartitions(t *testing.T) {
buffer := sink.NewSinkBuffer()
topic := "my-topic"
partition1 := kafka.TopicPartition{
Topic: &topic,
Partition: 1,
Offset: 1,
}
partition2 := kafka.TopicPartition{
Topic: &topic,
Partition: 2,
Offset: 1,
}
sinkMessage1 := sinkmodels.SinkMessage{
KafkaMessage: &kafka.Message{
TopicPartition: partition1,
},
}
sinkMessage2 := sinkmodels.SinkMessage{
KafkaMessage: &kafka.Message{
TopicPartition: partition2,
},
}
buffer.Add(sinkMessage1)
buffer.Add(sinkMessage2)
assert.Equal(t, 2, buffer.Size())
buffer.RemoveByPartitions([]kafka.TopicPartition{partition2})
assert.Equal(t, 1, buffer.Size())
}
|