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() // Unique identifier for each message (topic + partition + offset) 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 } // RemoveByPartitions removes messages from the buffer by partitions // Useful when partitions are revoked. 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) }