| package sink |
|
|
| import ( |
| "fmt" |
| "sync" |
|
|
| "github.com/confluentinc/confluent-kafka-go/v2/kafka" |
|
|
| sinkmodels "github.com/openmeterio/openmeter/openmeter/sink/models" |
| ) |
|
|
| type SinkBuffer struct { |
| mu sync.Mutex |
| data map[string]sinkmodels.SinkMessage |
| } |
|
|
| func NewSinkBuffer() *SinkBuffer { |
| return &SinkBuffer{ |
| data: map[string]sinkmodels.SinkMessage{}, |
| } |
| } |
|
|
| func (b *SinkBuffer) Size() int { |
| b.mu.Lock() |
| defer b.mu.Unlock() |
| return len(b.data) |
| } |
|
|
| func (b *SinkBuffer) Add(message sinkmodels.SinkMessage) { |
| b.mu.Lock() |
| defer b.mu.Unlock() |
| |
| key := message.KafkaMessage.String() |
| b.data[key] = message |
| } |
|
|
| type MessageTransformerFunc func(*sinkmodels.SinkMessage) |
|
|
| func (b *SinkBuffer) Dequeue(transformers ...MessageTransformerFunc) []sinkmodels.SinkMessage { |
| b.mu.Lock() |
| defer b.mu.Unlock() |
|
|
| messages := make([]sinkmodels.SinkMessage, 0, len(b.data)) |
|
|
| for key, message := range b.data { |
| for _, transformer := range transformers { |
| transformer(&message) |
| } |
|
|
| messages = append(messages, message) |
|
|
| delete(b.data, key) |
| } |
|
|
| return messages |
| } |
|
|
| |
| |
| func (b *SinkBuffer) RemoveByPartitions(partitions []kafka.TopicPartition) { |
| b.mu.Lock() |
| defer b.mu.Unlock() |
|
|
| partitionMap := map[string]bool{} |
| for _, topicPartition := range partitions { |
| key := topicPartitionKey(topicPartition) |
| partitionMap[key] = true |
| } |
|
|
| for key, message := range b.data { |
| topicKey := topicPartitionKey(message.KafkaMessage.TopicPartition) |
|
|
| if partitionMap[topicKey] { |
| delete(b.data, key) |
| } |
| } |
| } |
|
|
| func topicPartitionKey(partition kafka.TopicPartition) string { |
| var topic string |
| if partition.Topic != nil { |
| topic = *partition.Topic |
| } |
| return partitionKey(topic, partition.Partition) |
| } |
|
|
| func partitionKey(topic string, partition int32) string { |
| return fmt.Sprintf("%s-%d", topic, partition) |
| } |
|
|