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())
}