From 7de4e06120db786619a461adb52624bedecffb5e Mon Sep 17 00:00:00 2001 From: Krisztian Gacsal Date: Thu, 27 Aug 2026 15:56:25 +0200 Subject: [PATCH] refactor(ingest): use fixed Kafka event topic --- app/common/kafka.go | 16 +++ app/common/kafka_test.go | 42 +++++++ app/common/openmeter_server.go | 9 +- app/config/config_test.go | 1 + app/config/ingest.go | 2 + app/config/testdata/complete.yaml | 1 + cmd/server/wire_gen.go | 12 +- openmeter/ingest/kafkaingest/collector.go | 18 +-- .../ingest/kafkaingest/collector_test.go | 103 ++++++++++++++++++ 9 files changed, 182 insertions(+), 22 deletions(-) create mode 100644 app/common/kafka_test.go create mode 100644 openmeter/ingest/kafkaingest/collector_test.go diff --git a/app/common/kafka.go b/app/common/kafka.go index 8a72568cbc..ae9e7cea6f 100644 --- a/app/common/kafka.go +++ b/app/common/kafka.go @@ -23,6 +23,7 @@ var Kafka = wire.NewSet( ) var KafkaIngest = wire.NewSet( + NewEventTopic, NewKafkaIngestNamespaceHandler, ) @@ -31,6 +32,21 @@ var KafkaNamespaceResolver = wire.NewSet( wire.Bind(new(topicresolver.Resolver), new(*topicresolver.NamespacedTopicResolver)), ) +type EventTopic string + +// NewEventTopic returns the fixed Kafka destination for the ingest pipeline. +// Existing configurations fall back to the topic assigned to the default namespace. +func NewEventTopic( + ingestConfig config.KafkaIngestConfiguration, + namespaceConfig config.NamespaceConfiguration, +) EventTopic { + if ingestConfig.EventsTopic != "" { + return EventTopic(ingestConfig.EventsTopic) + } + + return EventTopic(fmt.Sprintf(ingestConfig.EventsTopicTemplate, namespaceConfig.Default)) +} + // TODO: add closer function? func NewKafkaProducer(conf config.KafkaIngestConfiguration, logger *slog.Logger, meta Metadata) (*kafka.Producer, error) { kafkaConfig := conf.CreateKafkaConfig() diff --git a/app/common/kafka_test.go b/app/common/kafka_test.go new file mode 100644 index 0000000000..b34ab3ff1b --- /dev/null +++ b/app/common/kafka_test.go @@ -0,0 +1,42 @@ +package common + +import ( + "testing" + + "github.com/stretchr/testify/assert" + + "github.com/openmeterio/openmeter/app/config" +) + +func TestNewEventTopic(t *testing.T) { + tests := []struct { + name string + ingestConfig config.KafkaIngestConfiguration + namespaceConfig config.NamespaceConfiguration + expected EventTopic + }{ + { + name: "explicit topic", + ingestConfig: config.KafkaIngestConfiguration{ + EventsTopic: "events", + EventsTopicTemplate: "om_%s_events", + }, + namespaceConfig: config.NamespaceConfiguration{Default: "default"}, + expected: EventTopic("events"), + }, + { + name: "default namespace fallback", + ingestConfig: config.KafkaIngestConfiguration{ + EventsTopicTemplate: "om_%s_events", + }, + namespaceConfig: config.NamespaceConfiguration{Default: "default"}, + expected: EventTopic("om_default_events"), + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + assert.Equal(t, test.expected, NewEventTopic(test.ingestConfig, test.namespaceConfig)) + }) + } +} diff --git a/app/common/openmeter_server.go b/app/common/openmeter_server.go index 9bea372ca9..f223619d0a 100644 --- a/app/common/openmeter_server.go +++ b/app/common/openmeter_server.go @@ -15,15 +15,14 @@ import ( "github.com/openmeterio/openmeter/openmeter/ingest/ingestadapter" "github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest" "github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest/serializer" - "github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest/topicresolver" watermillkafka "github.com/openmeterio/openmeter/openmeter/watermill/driver/kafka" pkgkafka "github.com/openmeterio/openmeter/pkg/kafka" ) func NewKafkaIngestCollector( - config config.KafkaIngestConfiguration, + ingestConfig config.KafkaIngestConfiguration, + eventTopic EventTopic, producer *kafka.Producer, - topicResolver topicresolver.Resolver, topicProvisioner pkgkafka.TopicProvisioner, logger *slog.Logger, tracer trace.Tracer, @@ -31,9 +30,9 @@ func NewKafkaIngestCollector( collector, err := kafkaingest.NewCollector( producer, serializer.NewJSONSerializer(), - topicResolver, + string(eventTopic), topicProvisioner, - config.Partitions, + ingestConfig.Partitions, logger, tracer, ) diff --git a/app/config/config_test.go b/app/config/config_test.go index 66896b8ae4..5a4bb53f83 100644 --- a/app/config/config_test.go +++ b/app/config/config_test.go @@ -122,6 +122,7 @@ func TestComplete(t *testing.T) { }, }, Partitions: 1, + EventsTopic: "om_explicit_events", EventsTopicTemplate: "om_%s_events", TopicProvisioner: TopicProvisionerConfig{ Enabled: true, diff --git a/app/config/ingest.go b/app/config/ingest.go index 410d18432e..bad6efab33 100644 --- a/app/config/ingest.go +++ b/app/config/ingest.go @@ -33,6 +33,7 @@ type KafkaIngestConfiguration struct { TopicProvisioner TopicProvisionerConfig Partitions int + EventsTopic string EventsTopicTemplate string // NamespaceDeletionEnabled defines whether deleting namespaces are allowed or not. @@ -177,6 +178,7 @@ func ConfigureIngestKafkaConfiguration(v *viper.Viper, prefixes ...string) { // Configure configures some defaults in the Viper instance. func ConfigureIngest(v *viper.Viper) { v.SetDefault("ingest.kafka.partitions", 1) + v.SetDefault("ingest.kafka.eventsTopic", "") v.SetDefault("ingest.kafka.eventsTopicTemplate", "om_%s_events") v.SetDefault("ingest.kafka.namespaceDeletionEnabled", false) diff --git a/app/config/testdata/complete.yaml b/app/config/testdata/complete.yaml index 94da95d976..4f8263ab67 100644 --- a/app/config/testdata/complete.yaml +++ b/app/config/testdata/complete.yaml @@ -40,6 +40,7 @@ ingest: saslUsername: user saslPassword: pass partitions: 1 + eventsTopic: om_explicit_events statsInterval: 5s brokerAddressFamily: any socketKeepAliveEnabled: true diff --git a/cmd/server/wire_gen.go b/cmd/server/wire_gen.go index fb82212b8c..610450c402 100644 --- a/cmd/server/wire_gen.go +++ b/cmd/server/wire_gen.go @@ -590,6 +590,7 @@ func initializeApplication(ctx context.Context, conf config.Configuration) (Appl return Application{}, nil, err } dedupeConfiguration := conf.Dedupe + eventTopic := common.NewEventTopic(kafkaIngestConfiguration, namespaceConfiguration) producer, err := common.NewKafkaProducer(kafkaIngestConfiguration, logger, commonMetadata) if err != nil { cleanup7() @@ -601,7 +602,7 @@ func initializeApplication(ctx context.Context, conf config.Configuration) (Appl cleanup() return Application{}, nil, err } - namespacedTopicResolver, err := common.NewNamespacedTopicResolver(kafkaIngestConfiguration) + collector, err := common.NewKafkaIngestCollector(kafkaIngestConfiguration, eventTopic, producer, topicProvisioner, logger, tracer) if err != nil { cleanup7() cleanup6() @@ -612,7 +613,7 @@ func initializeApplication(ctx context.Context, conf config.Configuration) (Appl cleanup() return Application{}, nil, err } - collector, err := common.NewKafkaIngestCollector(kafkaIngestConfiguration, producer, namespacedTopicResolver, topicProvisioner, logger, tracer) + ingestCollector, cleanup8, err := common.NewIngestCollector(dedupeConfiguration, collector, logger, meter, tracer) if err != nil { cleanup7() cleanup6() @@ -623,8 +624,9 @@ func initializeApplication(ctx context.Context, conf config.Configuration) (Appl cleanup() return Application{}, nil, err } - ingestCollector, cleanup8, err := common.NewIngestCollector(dedupeConfiguration, collector, logger, meter, tracer) + ingestService, err := common.NewIngestService(ingestCollector, logger, meter) if err != nil { + cleanup8() cleanup7() cleanup6() cleanup5() @@ -634,7 +636,7 @@ func initializeApplication(ctx context.Context, conf config.Configuration) (Appl cleanup() return Application{}, nil, err } - ingestService, err := common.NewIngestService(ingestCollector, logger, meter) + metrics, err := common.NewKafkaMetrics(meter) if err != nil { cleanup8() cleanup7() @@ -646,7 +648,7 @@ func initializeApplication(ctx context.Context, conf config.Configuration) (Appl cleanup() return Application{}, nil, err } - metrics, err := common.NewKafkaMetrics(meter) + namespacedTopicResolver, err := common.NewNamespacedTopicResolver(kafkaIngestConfiguration) if err != nil { cleanup8() cleanup7() diff --git a/openmeter/ingest/kafkaingest/collector.go b/openmeter/ingest/kafkaingest/collector.go index 286b53bc4c..b9218608f8 100644 --- a/openmeter/ingest/kafkaingest/collector.go +++ b/openmeter/ingest/kafkaingest/collector.go @@ -15,7 +15,6 @@ import ( "go.opentelemetry.io/otel/trace" "github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest/serializer" - "github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest/topicresolver" "github.com/openmeterio/openmeter/pkg/clock" pkgkafka "github.com/openmeterio/openmeter/pkg/kafka" kafkametrics "github.com/openmeterio/openmeter/pkg/kafka/metrics" @@ -40,7 +39,7 @@ func FromIngestedAt(s string) (time.Time, error) { type Collector struct { Producer *kafka.Producer Serializer serializer.Serializer - TopicResolver topicresolver.Resolver + Topic string TopicProvisioner pkgkafka.TopicProvisioner TopicPartitions int @@ -51,7 +50,7 @@ type Collector struct { func NewCollector( producer *kafka.Producer, serializer serializer.Serializer, - resolver topicresolver.Resolver, + topic string, provisioner pkgkafka.TopicProvisioner, partitions int, logger *slog.Logger, @@ -63,8 +62,8 @@ func NewCollector( if serializer == nil { return nil, fmt.Errorf("serializer is required") } - if resolver == nil { - return nil, fmt.Errorf("topic name resolver is required") + if topic == "" { + return nil, fmt.Errorf("topic is required") } if provisioner == nil { @@ -80,7 +79,7 @@ func NewCollector( return &Collector{ Producer: producer, Serializer: serializer, - TopicResolver: resolver, + Topic: topic, TopicProvisioner: provisioner, TopicPartitions: partitions, Logger: logger, @@ -105,12 +104,7 @@ func (s Collector) Ingest(ctx context.Context, namespace string, ev event.Event) span.End() }() - span.AddEvent("resolved namespace to kafka topic") - topicName, err := s.TopicResolver.Resolve(ctx, namespace) - if err != nil { - err = fmt.Errorf("failed to resolve namespace to topic name: %w", err) - return err - } + topicName := s.Topic span.SetAttributes(semconv.MessagingDestinationName(topicName)) // Make sure topic is provisioned diff --git a/openmeter/ingest/kafkaingest/collector_test.go b/openmeter/ingest/kafkaingest/collector_test.go new file mode 100644 index 0000000000..b35d74b348 --- /dev/null +++ b/openmeter/ingest/kafkaingest/collector_test.go @@ -0,0 +1,103 @@ +package kafkaingest + +import ( + "context" + "errors" + "testing" + + cloudevents "github.com/cloudevents/sdk-go/v2/event" + "github.com/confluentinc/confluent-kafka-go/v2/kafka" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/trace/noop" + + "github.com/openmeterio/openmeter/openmeter/testutils" + pkgkafka "github.com/openmeterio/openmeter/pkg/kafka" +) + +type recordingSerializer struct { + topics []string + err error +} + +func (s *recordingSerializer) SerializeKey(topic string, _ string, _ cloudevents.Event) ([]byte, error) { + s.topics = append(s.topics, topic) + + return nil, s.err +} + +func (s *recordingSerializer) SerializeValue(_ string, _ cloudevents.Event) ([]byte, error) { + return nil, nil +} + +func (s *recordingSerializer) GetFormat() string { + return "" +} + +func (s *recordingSerializer) GetKeySchemaId() int { + return 0 +} + +func (s *recordingSerializer) GetValueSchemaId() int { + return 0 +} + +type recordingTopicProvisioner struct { + topics []pkgkafka.TopicConfig +} + +func (p *recordingTopicProvisioner) Provision(_ context.Context, topics ...pkgkafka.TopicConfig) error { + p.topics = append(p.topics, topics...) + + return nil +} + +func (p *recordingTopicProvisioner) DeProvision(_ context.Context, _ ...string) error { + return nil +} + +func TestCollectorUsesFixedTopic(t *testing.T) { + serializerError := errors.New("stop before producing") + serializer := &recordingSerializer{err: serializerError} + provisioner := &recordingTopicProvisioner{} + + collector, err := NewCollector( + &kafka.Producer{}, + serializer, + "om_default_events", + provisioner, + 1, + testutils.NewDiscardLogger(t), + noop.NewTracerProvider().Tracer("test"), + ) + require.NoError(t, err) + + ev := cloudevents.New() + ev.SetID("event-id") + + for _, namespace := range []string{"default", "customer"} { + err = collector.Ingest(t.Context(), namespace, ev) + require.ErrorIs(t, err, serializerError) + } + + assert.Equal(t, []string{"om_default_events", "om_default_events"}, serializer.topics) + assert.Equal(t, []pkgkafka.TopicConfig{ + {Name: "om_default_events", Partitions: 1}, + {Name: "om_default_events", Partitions: 1}, + }, provisioner.topics) +} + +func TestNewCollectorRequiresTopic(t *testing.T) { + collector, err := NewCollector( + &kafka.Producer{}, + &recordingSerializer{}, + "", + &recordingTopicProvisioner{}, + 1, + testutils.NewDiscardLogger(t), + noop.NewTracerProvider().Tracer("test"), + ) + + require.EqualError(t, err, "topic is required") + assert.Nil(t, collector) +}